Menjaga sinkronisasi tabel Apache Iceberg target Lakehouse dengan tabel sumber Delta Lake 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 transaksi akan menimbulkan overhead operasional yang kompleks.
Dengan menggunakan pipeline pengambilan data perubahan (CDC) batch terjadwal Dataflow, Anda dapat terus menyinkronkan tabel sumber Delta Lake dengan tabel Apache Iceberg Lakehouse. Alih-alih memuat ulang tabel lengkap, tugas batch terjadwal mengekstrak pembaruan inkremental langsung dari log transaksi Delta Lake dan menerapkannya ke tabel target Anda sesuai jadwal yang dapat dikonfigurasi.
Pipeline CDC batch terjadwal menyediakan kemampuan berikut:
- Sinkronisasi inkremental yang hemat biaya: Secara berkala hanya memproses data yang berubah (operasi
INSERT,UPDATE, danDELETE), sehingga mengurangi overhead komputasi dan latensi data. - Pemetaan perubahan otomatis: Mem-parsing log Change Data Feed Delta Lake dan memetakan operasi perubahan ke tabel Iceberg Lakehouse target.
- Penjadwalan tugas yang fleksibel: Menjalankan tugas Dataflow mode batch sesuai jadwal (misalnya, interval 15 menit) yang sesuai dengan kebutuhan bisnis Anda.
Alur kerja migrasi berkelanjutan
Untuk menyiapkan migrasi berkelanjutan end-to-end dari Delta Lake ke Lakehouse, selesaikan urutan operasi berikut:
- Aktifkan properti
delta.enableChangeDataFeed(true) di tabel Delta Lake input. - Lakukan migrasi tabel awal ke ID atau stempel waktu commit setelah
mengaktifkan properti
delta.enableChangeDataFeed. Untuk mengetahui petunjuknya, lihat Mengimpor tabel Delta Lake ke Lakehouse menggunakan Dataflow. - Jalankan tugas CDC batch terjadwal untuk pembaruan data berikutnya dengan frekuensi (ID commit atau rentang stempel waktu) yang dipilih untuk persyaratan workload Anda.
Sebelum memulai
Untuk menyiapkan migrasi CDC batch terjadwal, pastikan Anda memiliki hal berikut:
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.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.
Tabel Delta Lake yang valid dan disimpan di bucket Cloud Storage. Direktori tabel harus berisi file data dan direktori log transaksi
_delta_log/. Selain itu, tabel harus memenuhi persyaratan berikut:- Properti
delta.enableChangeDataFeedharus diaktifkan (true) di tabel Delta Lake. - ID atau stempel waktu commit awal yang dikonfigurasi untuk pipeline CDC harus
berada setelah properti
delta.enableChangeDataFeeddiaktifkan.
- Properti
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 batch terjadwal dari Delta Lake 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 jika Anda memberikan daftar kolom kesetaraan (
equality_columns) yang valid. 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 Delta Lake sumber berubah, Anda harus memperbarui skema tabel Lakehouse target secara manual agar cocok sebelum menjalankan tugas sinkronisasi CDC.
- Mode eksekusi batch: CDC Delta Lake beroperasi dalam mode batch.
Pipeline menarik log transaksi sekali per eksekusi tugas antara batas commit awal dan akhir yang ditentukan (ID commit atau stempel waktu). Jika batas commit
akhir (
end_versionatauend_timestamp) tidak ditentukan, pipeline akan membaca hingga versi commit terbaru. - Penyimpanan sumber yang didukung: Data sumber harus berupa tabel Delta Lake yang valid dan disimpan di Cloud Storage (
gs://). Amazon S3 tidak didukung untuk sumber tabel Delta Lake. - Tabel berbasis katalog: Tabel Delta Lake berbasis katalog (misalnya, Unity Catalog) tidak didukung.
- Versi Delta Lake yang tidak didukung: Delta Lake versi 1.2.1 atau yang lebih lama tidak didukung karena Feed Data Perubahan (CDF) tidak didukung untuk versi ini.
- Feed Data Perubahan Otomatis: Feed Data Perubahan Otomatis tidak didukung karena tidak menulis file log perubahan CDF.
Membuat pipeline CDC secara terprogram
Untuk membuat dan menjalankan pipeline migrasi batch CDC Delta Lake secara terprogram, gunakan transformasi Managed.read di pipeline Apache Beam Anda.
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 menjalankan pipeline CDC batch dari Delta Lake 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.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 DeltaLakeCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// 1. Configure Delta Lake CDC reader
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "gs://BUCKET_NAME/TABLE_NAME");
// Starting commit bound (either start_version or start_timestamp must be provided):
readConfig.put("start_version", START_COMMIT_ID);
// Optional ending commit bound (end_version or end_timestamp, matching the start bound type):
// readConfig.put("end_version", END_COMMIT_ID);
readConfig.put("include_metadata_columns", Arrays.asList("_commit_version", "_change_type"));
// Read CDC events from Delta Lake
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromDeltaCDC",
Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig)).getSinglePCollection();
// 2. Configure BigLake REST Catalog properties for Lakehouse Iceberg
Map<String, String> catalogProps = new HashMap<>();
catalogProps.put("type", "rest");
catalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
catalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
catalogProps.put("header.x-goog-user-project", "PROJECT_ID");
catalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
catalogProps.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager");
// 3. Map Delta Lake change types to Iceberg CDC change types
Map<String, String> changeTypeMap = new HashMap<>();
changeTypeMap.put("insert", "INSERT");
changeTypeMap.put("delete", "DELETE");
changeTypeMap.put("update_preimage", "UPDATE_BEFORE");
changeTypeMap.put("update_postimage", "UPDATE_AFTER");
// 4. Configure Lakehouse Iceberg CDC writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_name", "lakehouse");
writeConfig.put("catalog_properties", catalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("sequence_number_column", "_commit_version");
writeConfig.put("equality_columns", Arrays.asList("PRIMARY_KEY_COLUMN"));
// change_type_column and change_type_map are only required for Portable Runner
writeConfig.put("change_type_column", "_change_type");
writeConfig.put("change_type_map", changeTypeMap);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
Ganti kode berikut:
BUCKET_NAME: nama bucket Cloud Storage yang berisi tabel Delta Lake.TABLE_NAME: nama direktori tabel Delta Lake sumber.START_COMMIT_ID: nomor versi commit awal untuk membaca peristiwa CDC Delta Lake.start_versionataustart_timestampharus diberikan. ID atau stempel waktu mulai commit harus setelah propertidelta.enableChangeDataFeeddiaktifkan di tabel sumber.END_COMMIT_ID: nomor versi commit akhir untuk membaca peristiwa CDC Delta Lake. Sebagai praktik terbaik, tentukan batas akhir (end_versionatauend_timestamp) yang cocok dengan jenis batas yang digunakan untuk awal (start_versionataustart_timestamp). Jika batas akhir tidak dicantumkan, pipeline akan membaca hingga versi commit terbaru.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.PRIMARY_KEY_COLUMN: nama kolom kunci primer yang digunakan untuk mengidentifikasi baris untuk update dan penghapusan CDC.
Periksa output tugas
Pastikan data CDC berhasil digabungkan ke tabel Lakehouse Anda:
Di konsol Google Cloud , buka halaman Studio BigQuery.
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.
Klik Run dan verifikasi hasilnya.
Langkah berikutnya
- Pelajari lebih lanjut Mengimpor tabel Delta Lake ke Lakehouse.
- Pelajari Penyerapan pengambilan data perubahan di Lakehouse.
- Pelajari lebih lanjut I/O Terkelola di Dataflow.