Bigtable HBase Beam 連接器

為協助您在 Dataflow 管道中使用 Bigtable,我們提供兩個開放原始碼的 Bigtable Beam I/O 連接器。

如果您要從 HBase 遷移至 Bigtable,或是應用程式會呼叫 HBase API,請使用本頁討論的 Bigtable HBase Beam 連接器 (CloudBigtableIO)。

在所有其他情況下,您都應搭配使用 Bigtable Beam 連接器 (BigtableIO) 和適用於 Java 的 Cloud Bigtable 用戶端,因為後者可與 Cloud Bigtable API 搭配運作。如要開始使用該連接器,請參閱「Bigtable Beam 連接器」。

如要進一步瞭解 Apache Beam 程式設計模型,請參閱 Beam 說明文件。

開始使用 HBase

Bigtable HBase Beam 連接器是以 Java 編寫而成,並以適用於 Java 的 Bigtable HBase 用戶端為基礎。這項服務與以 Apache Beam 為基礎的 Java 適用的 Dataflow SDK 2.x 相容。這個連結器的原始碼位於 GitHub 的 googleapis/java-bigtable-hbase 存放區。

本頁概要說明如何使用 Read 和 Write 轉換。

設定驗證方法

如要在本機開發環境中使用本頁的 Java 範例,請安裝並初始化 gcloud CLI,然後使用使用者憑證設定應用程式預設憑證。

  1. 安裝 Google Cloud CLI。

  2. 如果您使用外部識別資訊提供者 (IdP),請先 使用聯合身分登入 gcloud CLI。

  3. 如果您使用本機殼層,請為使用者帳戶建立本機驗證憑證:

    gcloud auth application-default login

    如果您使用 Cloud Shell,則不需要執行這項操作。

    如果系統傳回驗證錯誤,且您使用外部識別資訊提供者 (IdP),請確認您已 使用聯合身分登入 gcloud CLI。

詳情請參閱「 設定本機開發環境的驗證機制」。

如要瞭解如何為正式環境設定驗證機制,請參閱「 為在 Google Cloud上執行的程式碼設定應用程式預設憑證 」。

將連接器新增至 Maven 專案

如要將 Bigtable HBase Beam 連接器新增至 Maven 專案,請將 Maven 構件新增至 pom.xml 檔案做為依附元件:

<dependency>
  <groupId>com.google.cloud.bigtable</groupId>
  <artifactId>bigtable-hbase-beam</artifactId>
  <version>2.12.0</version>
</dependency>

指定 Bigtable 設定

建立選項介面,允許輸入內容來執行管道:

public interface BigtableOptions extends DataflowPipelineOptions {

  @Description("The Bigtable project ID, this can be different than your Dataflow project")
  @Default.String("bigtable-project")
  String getBigtableProjectId();

  void setBigtableProjectId(String bigtableProjectId);

  @Description("The Bigtable instance ID")
  @Default.String("bigtable-instance")
  String getBigtableInstanceId();

  void setBigtableInstanceId(String bigtableInstanceId);

  @Description("The Bigtable table ID in the instance.")
  @Default.String("mobile-time-series")
  String getBigtableTableId();

  void setBigtableTableId(String bigtableTableId);
}

從 Bigtable 讀取資料或將資料寫入 Bigtable 時,您必須提供 CloudBigtableConfiguration 設定物件。此物件指定您資料表的專案 ID 與執行個體 ID,以及資料表本身的名稱:

CloudBigtableTableConfiguration bigtableTableConfig =
    new CloudBigtableTableConfiguration.Builder()
        .withProjectId(options.getBigtableProjectId())
        .withInstanceId(options.getBigtableInstanceId())
        .withTableId(options.getBigtableTableId())
        .build();

如要讀取資料,請提供 CloudBigtableScanConfiguration 設定物件,以便指定 Apache HBase Scan 物件,限制及篩選讀取結果。詳情請參閱「從 Bigtable 讀取資料」一文。

從 Bigtable 讀取

如要從 Bigtable 資料表讀取資料,請將 Read 轉換套用至 CloudBigtableIO.read 作業的結果。Read 轉換會傳回 HBase Result 物件的 PCollection,其中 PCollection 中的每個元素都代表資料表中的單一資料列。

p.apply(Read.from(CloudBigtableIO.read(config)))
    .apply(
        ParDo.of(
            new DoFn<Result, Void>() {
              @ProcessElement
              public void processElement(@Element Result row, OutputReceiver<Void> out) {
                System.out.println(Bytes.toString(row.getRow()));
              }
            }));

根據預設,CloudBigtableIO.read 作業會傳回表格中的所有資料列。您可以使用 HBase Scan 物件,將讀取作業限制在資料表中的某個資料列鍵範圍內,或是對讀取結果套用篩選器。如要使用 Scan 物件,請將其納入 CloudBigtableScanConfiguration。

舉例來說,您可以新增 Scan,只傳回資料表中每列的第一個鍵/值組合,這在計算資料表中的列數時非常實用:

import com.google.cloud.bigtable.beam.CloudBigtableIO;
import com.google.cloud.bigtable.beam.CloudBigtableScanConfiguration;
import org.apache.beam.runners.dataflow.options.DataflowPipelineOptions;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.Read;
import org.apache.beam.sdk.options.Default;
import org.apache.beam.sdk.options.Description;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.hadoop.hbase.client.Result;
import org.apache.hadoop.hbase.client.Scan;
import org.apache.hadoop.hbase.filter.FirstKeyOnlyFilter;
import org.apache.hadoop.hbase.util.Bytes;

public class HelloWorldRead {
  public static void main(String[] args) {
    BigtableOptions options =
        PipelineOptionsFactory.fromArgs(args).withValidation().as(BigtableOptions.class);
    Pipeline p = Pipeline.create(options);

    Scan scan = new Scan();
    scan.setCacheBlocks(false);
    scan.setFilter(new FirstKeyOnlyFilter());

    CloudBigtableScanConfiguration config =
        new CloudBigtableScanConfiguration.Builder()
            .withProjectId(options.getBigtableProjectId())
            .withInstanceId(options.getBigtableInstanceId())
            .withTableId(options.getBigtableTableId())
            .withScan(scan)
            .build();

    p.apply(Read.from(CloudBigtableIO.read(config)))
        .apply(
            ParDo.of(
                new DoFn<Result, Void>() {
                  @ProcessElement
                  public void processElement(@Element Result row, OutputReceiver<Void> out) {
                    System.out.println(Bytes.toString(row.getRow()));
                  }
                }));

    p.run().waitUntilFinish();
  }

  public interface BigtableOptions extends DataflowPipelineOptions {
    @Description("The Bigtable project ID, this can be different than your Dataflow project")
    @Default.String("bigtable-project")
    String getBigtableProjectId();

    void setBigtableProjectId(String bigtableProjectId);

    @Description("The Bigtable instance ID")
    @Default.String("bigtable-instance")
    String getBigtableInstanceId();

    void setBigtableInstanceId(String bigtableInstanceId);

    @Description("The Bigtable table ID in the instance.")
    @Default.String("mobile-time-series")
    String getBigtableTableId();

    void setBigtableTableId(String bigtableTableId);
  }
}

寫入 Bigtable

如要寫入 Bigtable 資料表,請使用 apply a CloudBigtableIO.writeToTable 作業。您需要在 HBase Mutation 物件的 PCollection 上執行這項作業,這些物件可以包含 Put 和 Delete 物件。

Bigtable 資料表必須已存在,且必須定義適當的資料欄系列。Dataflow 連接器不會即時建立資料表和資料欄系列。您可以使用 cbt CLI 建立資料表並設定資料欄系列,也可以透過程式輔助方式完成這項操作。

如要寫入 Bigtable,請先建立 Dataflow 管道,以便透過網路序列化放置和刪除作業:

BigtableOptions options =
    PipelineOptionsFactory.fromArgs(args).withValidation().as(BigtableOptions.class);
Pipeline p = Pipeline.create(options);

一般來說,您需要執行轉換 (例如 ParDo),將輸出資料格式化為 HBase Put 或 Delete 物件的集合。以下範例顯示 DoFn 轉換,該轉換會取得目前值,並將其做為 Put 的資料列鍵。接著,您可以將 Put 物件寫入 Bigtable。

p.apply(Create.of("phone#4c410523#20190501", "phone#4c410523#20190502"))
    .apply(
        ParDo.of(
            new DoFn<String, Mutation>() {
              @ProcessElement
              public void processElement(@Element String rowkey, OutputReceiver<Mutation> out) {
                long timestamp = System.currentTimeMillis();
                Put row = new Put(Bytes.toBytes(rowkey));

                row.addColumn(
                    Bytes.toBytes("stats_summary"),
                    Bytes.toBytes("os_build"),
                    timestamp,
                    Bytes.toBytes("android"));
                out.output(row);
              }
            }))
    .apply(CloudBigtableIO.writeToTable(bigtableTableConfig));

如要啟用批次寫入流量控制,請將 BIGTABLE_ENABLE_BULK_MUTATION_FLOW_CONTROL 設為 true。這項功能會自動限制批次寫入要求的流量,並讓 Bigtable 自動調度資源功能自動新增或移除節點,以處理 Dataflow 工作。

CloudBigtableTableConfiguration bigtableTableConfig =
    new CloudBigtableTableConfiguration.Builder()
        .withProjectId(options.getBigtableProjectId())
        .withInstanceId(options.getBigtableInstanceId())
        .withTableId(options.getBigtableTableId())
        .withConfiguration(BigtableOptionsFactory.BIGTABLE_ENABLE_BULK_MUTATION_FLOW_CONTROL,
            "true")
        .build();
return bigtableTableConfig;

以下是完整的寫入範例,包括啟用批次寫入流量控制的變體。


import com.google.cloud.bigtable.beam.CloudBigtableIO;
import com.google.cloud.bigtable.beam.CloudBigtableTableConfiguration;
import com.google.cloud.bigtable.hbase.BigtableOptionsFactory;
import org.apache.beam.runners.dataflow.options.DataflowPipelineOptions;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.options.Default;
import org.apache.beam.sdk.options.Description;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.Create;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.hadoop.hbase.client.Mutation;
import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.util.Bytes;

public class HelloWorldWrite {

  public static void main(String[] args) {
    BigtableOptions options =
        PipelineOptionsFactory.fromArgs(args).withValidation().as(BigtableOptions.class);
    Pipeline p = Pipeline.create(options);

    CloudBigtableTableConfiguration bigtableTableConfig =
        new CloudBigtableTableConfiguration.Builder()
            .withProjectId(options.getBigtableProjectId())
            .withInstanceId(options.getBigtableInstanceId())
            .withTableId(options.getBigtableTableId())
            .build();

    p.apply(Create.of("phone#4c410523#20190501", "phone#4c410523#20190502"))
        .apply(
            ParDo.of(
                new DoFn<String, Mutation>() {
                  @ProcessElement
                  public void processElement(@Element String rowkey, OutputReceiver<Mutation> out) {
                    long timestamp = System.currentTimeMillis();
                    Put row = new Put(Bytes.toBytes(rowkey));

                    row.addColumn(
                        Bytes.toBytes("stats_summary"),
                        Bytes.toBytes("os_build"),
                        timestamp,
                        Bytes.toBytes("android"));
                    out.output(row);
                  }
                }))
        .apply(CloudBigtableIO.writeToTable(bigtableTableConfig));

    p.run().waitUntilFinish();
  }

  public interface BigtableOptions extends DataflowPipelineOptions {

    @Description("The Bigtable project ID, this can be different than your Dataflow project")
    @Default.String("bigtable-project")
    String getBigtableProjectId();

    void setBigtableProjectId(String bigtableProjectId);

    @Description("The Bigtable instance ID")
    @Default.String("bigtable-instance")
    String getBigtableInstanceId();

    void setBigtableInstanceId(String bigtableInstanceId);

    @Description("The Bigtable table ID in the instance.")
    @Default.String("mobile-time-series")
    String getBigtableTableId();

    void setBigtableTableId(String bigtableTableId);
  }

  public static CloudBigtableTableConfiguration batchWriteFlowControlExample(
      BigtableOptions options) {
    CloudBigtableTableConfiguration bigtableTableConfig =
        new CloudBigtableTableConfiguration.Builder()
            .withProjectId(options.getBigtableProjectId())
            .withInstanceId(options.getBigtableInstanceId())
            .withTableId(options.getBigtableTableId())
            .withConfiguration(BigtableOptionsFactory.BIGTABLE_ENABLE_BULK_MUTATION_FLOW_CONTROL,
                "true")
            .build();
    return bigtableTableConfig;
  }
}