Halaman ini menjelaskan karakteristik performa, per versi Apache Beam 2.75.0, untuk tugas streaming Dataflow yang membaca dari Apache Kafka dan menulis ke tabel Apache Iceberg. Halaman ini mengevaluasi perbedaan performa antara penulisan Apache Iceberg langsung dan penulisan yang dirutekan melalui Managed BigQuery API, serta membandingkan hasil ini dengan tolok ukur dasar dari pipeline Kafka ke BigQuery. Karena pengoptimalan untuk I/O Apache Iceberg sedang berlangsung, metrik performa ini dapat berubah.
Perbandingan tolok ukur tersedia di tiga konfigurasi pemetaan tanpa status utama (yang berarti membaca dari sumber, mengonversi pesan ke data, dan menulis ke sink tanpa melacak status atau menerapkan logika bisnis yang kompleks; disebut sebagai map_only atau mapping dalam tolok ukur):
- Kafka ke BigQuery (
map_only) (tolok ukur dari performa Kafka ke BigQuery) - Kafka ke Iceberg Langsung (
map_only,autosharding=false) - Kafka ke Iceberg menggunakan Managed BigQuery API (
map_only)
Selain itu, panduan ini mengevaluasi pola streaming Apache Iceberg langsung—seperti batching stateful menggunakan groupbykey—dan menjelaskan pertimbangan hilir penting terkait distribusi ukuran file, perilaku autosharding, dan latensi kueri sisi baca.
Metodologi pengujian
Tolok ukur dilakukan menggunakan resource berikut:
- Cluster Managed Service untuk Apache Kafka: Traffic dibuat menggunakan template Dataflow Streaming Data Generator.
- Throughput Input: 1 GBps
- Frekuensi Pesan: ~1.000.000 pesan per detik
- Format Pesan: Teks JSON dengan skema tetap (~1 KB per pesan)
- Partisi: 1.000 partisi Kafka
- Sink Tujuan:
- BigQuery: Tabel standar (tidak dipartisi) yang ditulis menggunakan BigQuery Storage Write API.
- Apache Iceberg: Katalog yang didukung oleh Cloud Storage. Sink Langsung dipartisi menggunakan
bucket(id, 64)(dikelompokkan ke dalam 64 shard pada kunci utama) dan menggunakan mode distribusihash.
Setelah penskalaan otomatis horizontal stabil, setiap konfigurasi pipeline berjalan dalam kondisi stabil selama 24 jam. Tolok ukur untuk setiap kasus pipeline dijalankan 3 kali secara terpisah, dan semua nilai yang dilaporkan mewakili rata-rata yang dihitung di seluruh proses tersebut untuk memastikan metrik performa yang berkelanjutan dan andal.
Performa penyerapan: Workload pemetaan
Pipeline pemetaan tanpa status membaca dari sumber, mengonversi format pesan ke data, dan menulis ke sink tanpa melacak status di seluruh data. Bagian berikut menganalisis arsitektur referensi yang berjalan pada 1 GBps.
Konfigurasi tugas
| Setelan | Kafka ke BigQuery (map_only) |
Kafka ke Iceberg Langsung (autosharding=false) |
Kafka ke Iceberg menggunakan Managed BigQuery API |
|---|---|---|---|
| Jenis mesin pekerja | e2-standard-2 |
e2-standard-4 |
e2-standard-4 |
| vCPU per Pekerja | 2 | 4 | 4 |
| RAM per Pekerja | 8 GB | 16 GB | 16 GB |
| Streaming Engine | Diaktifkan | Diaktifkan | Diaktifkan |
| Penskalaan Otomatis Horizontal | Diaktifkan | Diaktifkan | Diaktifkan |
| Frekuensi Pemicu | 5 detik | 60 detik | 60 detik |
Throughput dan penggunaan resource
Penulisan langsung ke file Parquet fisik dalam penyimpanan objek akan menimbulkan overhead I/O yang lebih tinggi daripada penyerapan streaming BigQuery. Dibandingkan dengan penulisan Iceberg langsung, penulisan perutean melalui Managed BigQuery API meningkatkan penggunaan CPU pekerja (~70% dibandingkan ~60%) dan mengurangi konsumsi Streaming Engine secara sederhana (~180 SECU/jam dibandingkan ~200 SECU/jam), meskipun persyaratan komputasi pekerja secara keseluruhan tetap serupa (~440 vCPU dibandingkan ~450 vCPU).
| Metrik | Kafka ke BigQuery (map_only) |
Kafka ke Iceberg Langsung (autosharding=false) |
Kafka ke Iceberg menggunakan Managed BigQuery API |
|---|---|---|---|
| Throughput Input Rata-Rata per Pekerja | ~15 MBps | ~9 MBps | ~9 MBps |
| Penggunaan CPU Rata-Rata | ~70% | ~60% | ~70% |
| Estimasi vCPU untuk Input 1 GBps | ~126 vCPU | ~450 vCPU | ~440 vCPU |
| Estimasi Pekerja untuk Input 1 GBps | ~63 pekerja | ~110 pekerja | ~110 pekerja |
| Estimasi SECU per Jam untuk 1 GBps | ~58 SECU/jam | ~200 SECU/jam | ~180 SECU/jam |
Profil latensi tulis
Penulisan Iceberg langsung menunjukkan latensi ekor (P99) yang parah karena batasan commit metadata penyimpanan objek. Penggunaan Managed BigQuery API akan menghilangkan lonjakan latensi ekor sekaligus mempertahankan latensi median yang rendah.
| Latensi Tulis End-to-End | Kafka ke BigQuery | Kafka ke Iceberg Langsung (autosharding=false) |
Kafka ke Iceberg menggunakan Managed BigQuery API |
|---|---|---|---|
| P50 (Median) | ~1.200 md | ~1.000 md | ~1.000 md |
| P95 | ~3.000 md | ~7.400 md | ~1.900 md |
| P99 (Ekor) | ~5.400 md | ~14.000 md | ~2.700 md |
Pertimbangan autosharding &pilihan desain
Bagian ini membahas implikasi autosharding pada ukuran file dan latensi pipeline saat menulis ke Apache Iceberg.
Alasan autosharding=false dipilih sebagai tolok ukur
Dalam pengujian awal, mengaktifkan autosharding menyebabkan ukuran file menyusut menjadi potongan kecil dan berfluktuasi secara arbitrer karena pemisahan shard dinamis yang dipicu oleh lonjakan beban tingkat thread lokal—bahkan di bawah beban input agregat yang konstan.
Untuk mempertahankan tata letak file Parquet yang stabil dan dapat diprediksi (~800 KB rata-rata) dan memastikan tolok ukur yang wajar tanpa flush prematur, autosharding=false dipilih untuk konfigurasi sink langsung.
Apa yang terjadi jika Anda menonaktifkan autosharding dibandingkan dengan tetap mengaktifkannya?
- Dengan
autosharding=false(Tolok Ukur): Anda akan mendapatkan ukuran file awal yang lebih besar (~800 KB rata-rata) dibandingkan dengan autosharding. Meskipun masih kecil dibandingkan dengan ukuran file Iceberg yang ideal (128–512 MB), hal ini memerlukan pemadatan hilir yang jauh lebih sedikit. Namun, komprominya adalah latensi ekor tulis yang tinggi (P99 mencapai ~14,0 detik) karena bottleneck metadata penyimpanan objek. - Jika autosharding diaktifkan: Dataflow akan menskalakan thread penulis secara dinamis untuk menyerap lonjakan throughput lokal, yang mengurangi latensi ekor tulis. Namun, hal ini akan mengganggu lapisan penyimpanan dengan menghasilkan file Parquet kecil dan terfragmentasi dalam volume tinggi (~100 KB atau lebih kecil). Ukuran file ini menunjukkan varians yang tinggi dan berfluktuasi secara arbitrer di seluruh proses (berkisar antara ~39 KB hingga ~100 KB rata-rata), sehingga meningkatkan kebutuhan akan pemeliharaan pemadatan hilir yang agresif.
Penyesuaian &Rekomendasi Partisi
Selama evaluasi, kami bereksperimen dengan berbagai nilai partisi tetap untuk tabel tujuan guna menemukan keseimbangan yang optimal. Kami menemukan bahwa penggunaan 64 bucket (misalnya, bucket(id, 64)) untuk partisi tabel tujuan menghasilkan ukuran file yang ditargetkan sekaligus mempertahankan penggunaan dan throughput yang layak. Pendekatan ini membantu kami mencocokkan manfaat performa autosharding sekaligus menghindari masalah fragmentasi ukuran file arbitrer yang terkait dengan penskalaan yang sepenuhnya dinamis.
Rekomendasi untuk Praktisi: Pelanggan dianjurkan untuk melakukan pengujian awal serupa dengan setelan partisi yang ditargetkan untuk menemukan titik optimal yang memaksimalkan paralelisme pipeline tanpa mengorbankan ukuran file Parquet.
Implikasi baca hilir: Ukuran file dan pemadatan
Meskipun metrik sisi tulis mendukung Managed BigQuery API untuk penyerapan Iceberg, efisiensi pipeline secara keseluruhan sangat bergantung pada performa baca hilir:
- Pembuatan File Kecil di Managed BigQuery API: Managed BigQuery API sering melakukan flush data untuk memastikan latensi tulis yang rendah. Perilaku ini menghasilkan file Parquet kecil dalam volume tinggi yang ditulis ke katalog Iceberg target.
- Dampak Latensi Kueri Baca: Mesin kueri (misalnya, Starburst/Trino, Apache Spark, BigQuery, Dremio) yang membaca tabel dengan jutaan file Parquet kecil akan menimbulkan overhead parsing metadata yang berat dan penalti pemindaian partisi.
- Persyaratan Pemadatan: Untuk mencegah penurunan performa baca saat menggunakan Managed BigQuery API (atau jika autosharding diaktifkan pada penulisan langsung), jalankan tugas pemeliharaan pemadatan Iceberg reguler (misalnya,
REWRITE DATA FILES). Overhead komputasi untuk pemadatan harus diperhitungkan dalam desain arsitektur secara keseluruhan. - Distribusi File Penulisan Langsung (
autosharding=false): Penulisan Iceberg langsung dengan sharding tetap menghasilkan file Parquet rata-rata yang lebih besar (~800 KB), sehingga menghasilkan tata letak yang kurang terfragmentasi untuk akses kueri langsung tanpa permintaan pemadatan langsung (meskipun masih di bawah rentang ideal).
Pipeline Iceberg langsung stateful (groupbykey)
Untuk mengevaluasi strategi batching manual, pengelompokan kunci stateful (groupbykey) diuji terhadap pipeline Kafka ke Iceberg Langsung (map_only, autosharding=false) tolok ukur. Kedua konfigurasi menulis file Parquet langsung ke penyimpanan objek.
Perbandingan tolok ukur
| Metrik / Fitur | Tolok Ukur Sink Langsung (autosharding=false) |
Sink Langsung Stateful (groupbykey) |
Dampak Performa |
|---|---|---|---|
| Estimasi vCPU untuk 1 GBps | ~450 vCPU | ~520 vCPU | ~+16% komputasi diperlukan |
| Penggunaan CPU Rata-Rata | ~60% | ~50% | ~-17% efisiensi pekerja |
| Estimasi SECU/jam untuk 1 GBps | ~200 SECU/jam | ~300 SECU/jam | ~+50% beban Streaming Engine |
| Ukuran File Rata-Rata | ~800 KB | ~100 KB | Menghasilkan batch file yang lebih kecil |
| Latensi P50 | ~1.000 md | ~1.200 md | ~+20% median lebih lambat |
| Latensi P95 | ~7.400 md | ~5.500 md | ~-26% latensi lebih rendah |
| Latensi P99 | ~14.000 md | ~13.000 md | Perubahan latensi ekor marginal |
Analisis kompromi
- Overhead Streaming Engine: Menambahkan langkah
groupbykeystateful mengharuskan Beam menyimpan status perantara di seluruh batas jendela. Hal ini meningkatkan konsumsi Unit Komputasi Streaming Engine sebesar ~50% (dari ~200 SECU/jam menjadi ~300 SECU/jam). - Latensi Buffering: Agregasi kunci manual memperkenalkan buffering jendela wajib, yang meningkatkan latensi tulis median (P50) menjadi ~1.200 md dan latensi P95 menjadi ~5,5 detik.
Pipeline terbalik: Streaming dari Iceberg ke Kafka
Untuk mengevaluasi kemampuan lakehouse dua arah, tolok ukur juga dilakukan untuk data streaming yang mengalir terbalik—membaca aliran hanya-tambah dari tabel Apache Iceberg dan memublikasikan kembali ke Apache Kafka.
Konfigurasi dan efisiensi tugas
Tidak seperti pipeline penyerapan yang harus menangani penulisan file penyimpanan objek yang berat atau bottleneck commit metadata, membaca dan melakukan streaming perubahan dari Iceberg beroperasi dengan efisiensi tinggi:
| Metrik | Iceberg ke Kafka (Hanya-Tambah, Tepat Sekali) |
|---|---|
| Jenis mesin pekerja | e2-standard-4 |
| Estimasi vCPU untuk Input 1 GBps | ~30 vCPU |
| Estimasi Pekerja untuk Input 1 GBps | ~7 pekerja |
| Estimasi SECU per Jam untuk 1 GBps | ~0,2 SECU/jam |
Poin-poin penting untuk pipeline terbalik
- Overhead Komputasi yang Jauh Lebih Rendah: Membaca dan memproyeksikan aliran CDC dari Iceberg memerlukan resource komputasi yang jauh lebih sedikit (~30 vCPU dibandingkan ~450 vCPU untuk penulisan langsung) karena menghindari tugas berat partisi, encoding, dan commit file Parquet dalam volume besar ke penyimpanan objek.
- Efisiensi Resource: Konsumsi atau replikasi berbasis peristiwa hilir dari format lakehouse kembali ke lapisan streaming sangat efisien dibandingkan dengan jalur penyerapan masuk.
Ringkasan rekomendasi arsitektur
| Pola Arsitektur | Latensi Tulis P99 | Tata Letak File | Pertimbangan Baca Hilir |
|---|---|---|---|
Kafka ke BigQuery (map_only) |
~5,4 detik | T/A | Optimal (Managed BigQuery Storage Engine) |
| Kafka ke Iceberg menggunakan Managed BigQuery API | ~2,7 detik | File yang Sangat Kecil secara Arbitrer | Memerlukan pemadatan berkala untuk pembacaan bervolume tinggi |
Kafka ke Iceberg Langsung (autosharding=false) |
~14,0 detik | ~800 KB | Baik (Ukuran file awal yang lebih besar, permintaan pemadatan yang lebih rendah) |
Kafka ke Iceberg Langsung (groupbykey) |
~13,0 detik | ~100 KB | Sedang (Overhead komputasi dan status yang lebih tinggi) |
Perkirakan biaya
Anda dapat memperkirakan biaya dasar pipeline Anda sendiri yang sebanding dengan penagihan Berbasis Resource menggunakan Google Cloud kalkulator harga, sebagai berikut:
- Buka kalkulator harga.
- Klik Tambahkan ke estimasi.
- Pilih Dataflow.
- Untuk Jenis layanan, pilih "Dataflow Classic".
- Pilih Setelan lanjutan untuk menampilkan kumpulan opsi lengkap.
- Pilih lokasi tempat tugas berjalan.
- Untuk Jenis pekerjaan, pilih "Streaming".
- Pilih Aktifkan Streaming Engine.
- Masukkan informasi untuk jam operasi tugas, node pekerja, mesin pekerja, dan penyimpanan Persistent Disk.
- Masukkan perkiraan jumlah Unit Komputasi Streaming Engine.
Penggunaan resource dan skala biaya kira-kira linear dengan throughput input, meskipun untuk tugas kecil dengan hanya beberapa pekerja, total biaya didominasi oleh biaya tetap. Sebagai titik awal, Anda dapat mengekstrapolasi jumlah node pekerja dan konsumsi resource dari hasil tolok ukur.
Misalnya, Anda menjalankan pipeline menggunakan arsitektur Kafka ke Iceberg Langsung (autosharding=false), dengan kecepatan data input 100 MBps. Berdasarkan hasil tolok ukur untuk pipeline 1 GBps, Anda dapat memperkirakan persyaratan resource sebagai berikut:
- Faktor Penskalaan: (100 MBps) / (1024 MBps) = ~0,1
- Node pekerja yang diproyeksikan: 110 pekerja × 0,1 = ~11 pekerja
- Jumlah Unit Komputasi Streaming Engine yang diproyeksikan per jam: 200 × 0,1 = ~20 unit per jam
Nilai ini hanya boleh digunakan sebagai perkiraan awal. Throughput dan biaya sebenarnya dapat sangat bervariasi, berdasarkan faktor-faktor seperti jenis mesin, distribusi ukuran pesan, kode pengguna, jenis agregasi, paralelisme kunci, dan ukuran jendela. Untuk mengetahui informasi selengkapnya, lihat Praktik terbaik untuk pengoptimalan biaya Dataflow.
Menjalankan pipeline pengujian
Untuk men-deploy tugas streaming Apache Iceberg menggunakan template Dataflow Flex, gunakan
gcloud dataflow flex-template run
perintah.
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--region=REGION \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
--enable-streaming-engine \
--parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'
Ganti kode berikut:
JOB_NAME: nama tugas Dataflow AndaPROJECT_ID: ID proyek Google Cloud AndaREGION: the Google Cloud region tempat tugas Anda berjalan (misalnya,us-central1)KAFKA_BOOTSTRAP_ADDRESS: alamat bootstrap cluster Apache Kafka AndaKAFKA_TOPIC: nama topik Kafka AndaICEBERG_TABLE_IDENTIFIER: ID tabel Iceberg target AndaCATALOG_NAME: nama katalog Iceberg AndaCATALOG_TYPE: jenis katalog yang akan digunakan (misalnya,hadoopataubigquery)BUCKET_NAME: nama bucket Cloud Storage untuk lokasi gudang AndaSCHEMA_DEFINITION: definisi skema untuk data topik Kafka Anda (misalnya,{"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})