Memigrasikan data pengambilan data perubahan (CDC) Apache Iceberg ke Lakehouse menggunakan Dataflow

Menjaga sinkronisasi tabel Apache Iceberg target Lakehouse dengan tabel sumber Apache Iceberg yang sering diperbarui bisa jadi sulit. Melakukan pemuatan ulang tabel penuh untuk merekam perubahan set data akan meningkatkan biaya komputasi dan latensi pemrosesan, serta membangun pipeline kustom untuk memproses log pengambilan data perubahan (CDC) akan menimbulkan overhead operasional yang kompleks.

Dengan menggunakan pipeline CDC batch atau streaming Dataflow, Anda dapat menyinkronkan tabel sumber Apache Iceberg secara terus-menerus atau berkala dengan tabel Apache Iceberg Lakehouse. Alih-alih memuat ulang tabel lengkap, pipeline CDC mengekstrak update inkremental langsung dari log perubahan Apache Iceberg sumber Anda dan menerapkannya ke tabel target Anda.

Pipeline CDC menyediakan kemampuan berikut:

  • Sinkronisasi inkremental yang hemat biaya: Hanya memproses data yang berubah (operasi INSERT, UPDATE, dan DELETE), sehingga mengurangi overhead komputasi dan latensi data.
  • Mode eksekusi yang fleksibel: Mendukung mode batch dan streaming agar sesuai dengan persyaratan bisnis Anda.

Pipeline CDC mendukung dua mode eksekusi:

  • Mode CDC batch: Membaca pembaruan inkremental secara berkala dari tabel Apache Iceberg sumber Anda di antara batas snapshot awal dan akhir yang ditentukan.
  • Mode CDC streaming: Terus-menerus melakukan polling pada tabel Apache Iceberg sumber Anda untuk mendapatkan commit perubahan baru dengan frekuensi yang dapat dikonfigurasi (default 1 menit, minimum 1 detik) dan menerapkannya ke tabel Lakehouse Anda dalam waktu mendekati real-time.

Sebelum memulai

Untuk menyiapkan migrasi CDC dari Apache Iceberg ke Lakehouse, pastikan Anda memiliki hal berikut:

  1. Aktifkan Dataflow API, BigQuery API, dan Lakehouse API, jika ada yang belum diaktifkan.

    Peran yang diperlukan untuk mengaktifkan API

    Untuk mengaktifkan API, Anda memerlukan izin serviceusage.services.enable. Jika Anda membuat project, kemungkinan Anda sudah memiliki izin ini melalui peran Pemilik (roles/owner). Jika tidak, Anda bisa mendapatkan izin ini melalui peran Admin Penggunaan Layanan (roles/serviceusage.serviceUsageAdmin). Pelajari cara memberikan peran.

    Aktifkan API

  2. Untuk mendapatkan izin yang diperlukan guna membuat resource, minta administrator Anda untuk memberi Anda peran Identity and Access Management (IAM) yang diperlukan di project Anda.

  3. Tabel Apache Iceberg sumber yang valid dan disimpan di Cloud Storage atau terdaftar di katalog atau metastore REST Iceberg.

  4. Katalog Iceberg Lakehouse yang ada untuk menerima data yang disinkronkan.

    Jika namespace target tidak ada, pipeline CDC akan otomatis membuatnya. Jika tabel Lakehouse target tidak ada, pipeline CDC akan otomatis membuatnya jika Anda memberikan daftar kolom persamaan yang valid.

Dukungan dan batasan

Migrasi CDC dari Apache Iceberg ke Lakehouse memiliki pertimbangan berikut:

  • Persyaratan SDK: Memerlukan Apache Beam SDK versi 2.77.0 atau yang lebih baru.
  • Pembuatan tabel tujuan: Pembuatan tabel Lakehouse tujuan didukung. Jika tabel target tidak ada, pipeline CDC akan otomatis membuatnya, asalkan Anda memberikan daftar kolom kesetaraan yang valid (equality_columns). Hanya katalog Lakehouse yang ada yang diperlukan; pipeline CDC akan otomatis membuat namespace jika tidak ada.
  • Evolusi skema: Perubahan skema otomatis tidak didukung. Jika skema tabel Apache Iceberg sumber berubah, Anda harus memperbarui skema tabel Lakehouse target secara manual agar cocok sebelum menjalankan tugas sinkronisasi CDC.
  • Mode eksekusi:
    • CDC Batch: Menarik perubahan log transaksi sekali per eksekusi tugas antara batas snapshot awal dan akhir yang ditentukan (ID snapshot atau stempel waktu).
    • Streaming CDC: Terus-menerus melakukan polling pada tabel Apache Iceberg sumber untuk mendapatkan log commit baru. Frekuensi polling dapat dikonfigurasi (defaultnya 1 menit, minimum 1 detik).

Membuat pipeline CDC menggunakan Java

Untuk membuat dan menjalankan pipeline migrasi batch atau streaming CDC Apache Iceberg secara terprogram menggunakan Apache Beam Java SDK, gunakan transformasi Managed.read dengan Managed.ICEBERG_CDC.

Contoh Java berikut menunjukkan cara mengonfigurasi dan mengeksekusi pipeline CDC Apache Iceberg untuk workload batch dan streaming.

Pipeline CDC batch

Contoh Java berikut mengonfigurasi pipeline batch yang membaca peristiwa CDC Apache Iceberg dalam batas cuplikan awal dan akhir yang ditentukan:

Java

Tambahkan dependensi berikut ke file pom.xml Anda:

<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-managed</artifactId>
  <version>${beam.version}</version>
</dependency>
<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-io-iceberg</artifactId>
  <version>${beam.version}</version>
</dependency>

Cuplikan Java berikut menunjukkan cara mengonfigurasi dan mengeksekusi pipeline CDC batch dari tabel Iceberg ke Lakehouse:

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;

public class IcebergBatchCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // Configure source catalog properties
    Map<String, String> sourceCatalogProps = new HashMap<>();
    sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
    sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
    // Add additional properties if required by your catalog type (for example, uri)

    // Configure batch Iceberg CDC reader with snapshot bounds
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
    readConfig.put("catalog_properties", sourceCatalogProps);
    readConfig.put("start_snapshot_id", START_SNAPSHOT_ID);
    readConfig.put("end_snapshot_id", END_SNAPSHOT_ID);
    readConfig.put(
        "include_metadata_columns",
        Arrays.asList(
            IcebergCdcMetadataColumns.CHANGE_TYPE,
            IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));

    // Read batch CDC events from source Iceberg table
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromIcebergCDC",
            Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();

    // Configure destination Lakehouse REST Catalog properties for Lakehouse Iceberg
    Map<String, String> destinationCatalogProps = new HashMap<>();
    destinationCatalogProps.put("type", "rest");
    destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
    destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
    destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
    destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
    destinationCatalogProps.put("rest.auth.type", "google");

    // Configure Lakehouse Iceberg writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_properties", destinationCatalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);

    // Write CDC rows to target Lakehouse Iceberg table
    cdcRows.apply(
        "WriteToLakehouse",
        Managed.write(Managed.ICEBERG).withConfig(writeConfig));

    pipeline.run();
  }
}

Ganti kode berikut:

  • SOURCE_CATALOG_TYPE: jenis katalog Iceberg sumber (misalnya, hadoop atau hive).
  • SOURCE_WAREHOUSE_LOCATION: lokasi gudang katalog Iceberg sumber (misalnya, s3://source-warehouse atau gs://source-warehouse).
  • SOURCE_NAMESPACE: namespace tabel Iceberg sumber.
  • SOURCE_TABLE: nama tabel Iceberg sumber.
  • START_SNAPSHOT_ID: ID snapshot awal untuk pembacaan CDC batch.
  • END_SNAPSHOT_ID: ID snapshot akhir untuk pembacaan CDC batch.
  • WAREHOUSE_BUCKET: nama bucket Cloud Storage yang digunakan sebagai gudang katalog Lakehouse.
  • PROJECT_ID: Google Cloud Project ID Anda.
  • TARGET_NAMESPACE: namespace tabel Lakehouse target.
  • TARGET_TABLE: nama tabel Lakehouse target.

Pipeline CDC streaming

Contoh Java berikut mengonfigurasi pipeline streaming yang terus-menerus meminta polling pada tabel Apache Iceberg sumber untuk perubahan CDC baru:

Java

Tambahkan dependensi berikut ke file pom.xml Anda:

<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-managed</artifactId>
  <version>${beam.version}</version>
</dependency>
<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-io-iceberg</artifactId>
  <version>${beam.version}</version>
</dependency>

Cuplikan Java berikut menunjukkan cara mengonfigurasi dan mengeksekusi pipeline CDC streaming dari tabel Iceberg ke Lakehouse:

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;

public class IcebergStreamingCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // Configure source catalog properties
    Map<String, String> sourceCatalogProps = new HashMap<>();
    sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
    sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
    // Add additional properties if required by your catalog type (for example, uri)

    // Configure streaming Iceberg CDC reader with continuous polling
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
    readConfig.put("catalog_properties", sourceCatalogProps);
    readConfig.put("streaming", true);
    readConfig.put("poll_interval_seconds", POLL_INTERVAL_SECONDS);
    readConfig.put(
        "include_metadata_columns",
        Arrays.asList(
            IcebergCdcMetadataColumns.CHANGE_TYPE,
            IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));

    // Read streaming CDC events from source Iceberg table
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromIcebergCDC",
            Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();

    // Configure destination BigLake REST Catalog properties for Lakehouse Iceberg
    Map<String, String> destinationCatalogProps = new HashMap<>();
    destinationCatalogProps.put("type", "rest");
    destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
    destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
    destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
    destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
    destinationCatalogProps.put("rest.auth.type", "google");

    // Configure Lakehouse Iceberg writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_properties", destinationCatalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);

    // Write CDC rows to target Lakehouse Iceberg table
    cdcRows.apply(
        "WriteToLakehouse",
        Managed.write(Managed.ICEBERG).withConfig(writeConfig));

    pipeline.run();
  }
}

Ganti kode berikut:

  • SOURCE_CATALOG_TYPE: jenis katalog Iceberg sumber (misalnya, hadoop atau hive).
  • SOURCE_WAREHOUSE_LOCATION: lokasi gudang katalog Iceberg sumber (misalnya, s3://source-warehouse atau gs://source-warehouse).
  • SOURCE_NAMESPACE: namespace tabel Iceberg sumber.
  • SOURCE_TABLE: nama tabel Iceberg sumber.
  • POLL_INTERVAL_SECONDS: frekuensi polling berkelanjutan dalam detik (default 60, minimum 1).
  • WAREHOUSE_BUCKET: nama bucket Cloud Storage yang digunakan sebagai gudang katalog Lakehouse.
  • PROJECT_ID: Google Cloud Project ID Anda.
  • TARGET_NAMESPACE: namespace tabel Lakehouse target.
  • TARGET_TABLE: nama tabel Lakehouse target.

Periksa output tugas

Pastikan data CDC berhasil digabungkan ke tabel Lakehouse Anda:

  1. Di konsol Google Cloud , buka halaman Studio BigQuery.

    Buka BigQuery

  2. Di editor kueri, jalankan kueri SQL untuk memverifikasi data yang disinkronkan:

    SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;
    

    Ganti kode berikut:

    • PROJECT_ID: Google Cloud Project ID Anda.
    • CATALOG: nama katalog Lakehouse Anda.
    • NAMESPACE: namespace tabel Lakehouse Anda.
    • TABLE_NAME: nama tabel Lakehouse target Anda.
  3. Klik Run dan verifikasi hasilnya.

Langkah berikutnya