Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Important
Fitur ini ada di Pratinjau Umum.
Gunakan Streaming Terstruktur untuk menulis ke Lakebase dengan batching bawaan, percobaan ulang otomatis, dan autentikasi yang dikelola ruang kerja.
Kapan menggunakan sink Lakebase
Gunakan sink Lakebase untuk penulisan streaming berlatensi rendah ke Lakebase. Sink ini tidak mengharuskan Anda menerapkan fungsi kustom foreachBatch untuk menangani batching, manajemen koneksi, dan penanganan kesalahan.
Kasus penggunaan umum meliputi:
- Perbarui database aplikasi secara real time untuk dasbor operasional atau fitur yang menghadap pelanggan.
- Sinkronkan data yang terus berubah, seperti hasil streaming yang diagregasi atau difilter, ke dalam basis data transaksional.
- Tulis output kueri Streaming Terstruktur ke dalam tabel Lakebase dengan latensi sub-detik menggunakan mode real time.
Untuk menyinkronkan data dari tabel Lakebase ke Delta Lake di Lakehouse, arah sebaliknya, lihat Umpan Data Perubahan Lakebase.
Persyaratan
- Databricks Runtime 18 ke atas
- Komputasi klasik dengan mode akses khusus atau standar.
- Database Lakebase
Sambungkan ke database
Sink Lakebase mendukung metode koneksi berikut:
Tabel Lakebase terdaftar di Unity Catalog
Untuk tabel Lakebase yang terdaftar di Unity Catalog, konektor secara otomatis mengelola kredensial dan menggunakan identitas pengguna atau perwakilan layanan yang menjalankan kueri. Jika tabel tidak ada, konektor akan membuat tabel.
Untuk mendaftarkan database Lakebase dengan Unity Catalog, lihat Mendaftarkan database Lakebase di Unity Catalog.
Untuk menulis ke tabel Lakebase, gunakan metode .toTable() dengan nama tabel yang memenuhi syarat secara lengkap, catalog.schema.table. Contoh berikut menunjukkan opsi yang diperlukan, ditambah opsi opsional upsertkey :
Python
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
Scala
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
Ganti placeholder berikut:
-
<catalog>.<schema>.<table>: Nama lengkap yang memenuhi syarat dari tabel target.catalogadalah katalog Unity Catalog yang Anda buat saat mendaftarkan database Lakebase, lihat Mendaftarkan database Lakebase di Katalog Unity. Jika tabel tidak ada, konektor akan membuatnya. -
<primary-key-column>: Opsional. Daftar kolom yang dipisahkan koma yang membentuk kunci upsert, misalnyaidatauuser_id,event_type. Jika Anda menghilangkanupsertkey, sink menyimpulkan kunci dari kunci utama tabel target. Lihat Perilaku upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Jalur volume Unity Catalog tempat kueri menyimpan checkpoint-nya. Anda juga dapat menggunakan URI penyimpanan objek cloud. Lokasi harus berupa penyimpanan yang dapat Anda tulis, bukan disk lokal, dan harus unik untuk setiap kueri streaming. Ini tidak bergantung pada tabel target. Lihat checkpoint Streaming Terstruktur.
Untuk konfigurasi opsional, seperti batchsize dan batchinterval, lihat Opsi konfigurasi.
Tabel Lakebase tidak terdaftar di Unity Catalog
Untuk tabel Lakebase yang tidak terdaftar di Unity Catalog, konektor secara otomatis mengelola kredensial dan menggunakan identitas pengguna atau perwakilan layanan yang menjalankan kueri. Jika tabel tidak ada, konektor akan membuat tabel.
Untuk menulis ke tabel Lakebase, gunakan opsi endpoint dan dbtable. Contoh berikut juga mencakup opsi opsional database dan upsertkey:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") # Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") // Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
Ganti placeholder berikut:
-
<project-id>.<branch-id>.<endpoint-id>: Titik akhir Lakebase Anda. Temukan ketiga nilai dalam Nama sumber daya pada menu Dapatkan ID dari tab Komputasi , yang memiliki formatprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Lihat Pengidentifikasi komputasi. -
<database>: Opsional. Nama database Postgres yang menjadi target. Secara default menjadidatabricks_postgres. Lihat Mengelola database. -
<schema>.<table>: Tabel target dalam formatschema.table. Jika Anda menghilangkan skema, sink menggunakan skemapublic. Gunakan pengidentifikasi sederhana yang dimulai dengan huruf atau garis bawah dan hanya berisi huruf, angka, dan garis bawah; pengidentifikasi yang dikutip dan karakter khusus, seperti tanda hubung, tidak didukung. -
<primary-key-column>: Opsional. Daftar kolom yang dipisahkan koma yang membentuk kunci upsert, misalnyaidatauuser_id,event_type. Jika Anda menghilangkanupsertkey, sink menyimpulkan kunci dari kunci utama tabel target. Lihat Perilaku upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Jalur volume Unity Catalog tempat kueri menyimpan checkpoint-nya. Anda juga dapat menggunakan URI penyimpanan objek cloud. Lokasi harus berupa penyimpanan yang dapat Anda tulis, bukan disk lokal, dan harus unik untuk setiap kueri streaming. Ini tidak bergantung pada tabel target. Lihat checkpoint Streaming Terstruktur.
Untuk konfigurasi opsional, seperti batchsize dan batchinterval, lihat Opsi konfigurasi.
Opsi konfigurasi
Sink menghasilkan error jika ada opsi yang tidak dikenali, JDBC_STREAMING_SINK_INVALID_OPTIONS.
Opsi berikut berlaku untuk semua metode koneksi:
| Key | Default | Description |
|---|---|---|
batchinterval |
100 milliseconds |
Fakultatif. Waktu maksimum untuk menyimpan baris di buffer sebelum dikosongkan. Contohnya, "50 milliseconds". |
batchsize |
1000 |
Fakultatif. Jumlah maksimum baris untuk setiap transaksi database. |
checkpointLocation |
Tidak | Required. Jalur ke direktori titik pemeriksaan, seperti volume Katalog Unity (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Harus unik untuk setiap kueri. Lihat checkpoint Streaming Terstruktur. |
upsertkey |
Tidak | Fakultatif. Daftar nama kolom yang dipisahkan koma yang membentuk kunci upsert. Misalnya, "id" atau "user_id,event_type". Jika Anda menentukan upsertkey, kolom harus cocok dengan kunci utama tabel, atau kueri gagal. Jika Anda menghilangkannya, sink akan menggunakan kunci utama secara otomatis. Untuk informasi selengkapnya, lihat Perilaku upsert. |
Tabel Lakebase tidak terdaftar di Unity Catalog
Opsi berikut berlaku saat Anda menyambungkan ke tabel Lakebase yang tidak terdaftar di Katalog Unity:
| Key | Default | Description |
|---|---|---|
database |
databricks_postgres |
Fakultatif. Nama database PostgreSQL target. |
dbtable |
Tidak | Required. Nama tabel target dalam schema.table format. Jika Anda tidak menentukan skema, nilai skema defaultnya adalah public. Gunakan pengidentifikasi sederhana yang dimulai dengan huruf atau garis bawah dan hanya berisi huruf, angka, dan garis bawah. Jangan gunakan tanda kutip pada nama tabel atau skema; pengidentifikasi yang diberi tanda kutip dan nama yang mengandung karakter khusus, seperti tanda hubung, tidak didukung. |
endpoint |
Tidak | Required. Endpoint Lakebase, dalam format project_id.branch_id atau project_id.branch_id.endpoint_id.
endpoint_id bersifat opsional; jika Anda menghilangkannya dan cabang memiliki satu titik akhir baca-tulis, sink memilih titik akhir tersebut secara default. |
Perilaku upsert
Ketika kunci upsert tersedia, baik ditetapkan melalui upsertkey maupun disimpulkan oleh sink dari kunci primer tabel, sink melakukan upsert ke tabel dengan sintaks INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ... PostgreSQL.
Ketika tidak ada kunci upsert, sink melakukan penyisipan. Mode keluaran dari sebuah kueri tidak memengaruhi operasi upsert atau insert.
Kolom upsertkey harus:
- Merupakan subset tidak kosong dari kolom-kolom DataFrame.
- Cocokkan
PRIMARY KEYdengan tabel target secara tepat. Jika kolom yang Anda tentukan tidak cocok dengan kunci utama, kueri akan gagal. - Harus berupa tipe yang dapat dibandingkan, seperti tipe numerik atau string. Untuk mencegah kebuntuan database selama penulisan bersamaan, sink mengurutkan baris menurut kunci upsert dalam setiap batch. Kunci upsert tidak mendukung tipe kompleks atau tipe struct.
Nama kolom secara otomatis diberi tanda kutip sesuai default PostgreSQL, tanda kutip ganda ", yang menangani kata kunci yang dicadangkan dan nama dengan huruf besar-kecil campuran.
Nama tabel dan skema harus menggunakan pengidentifikasi sederhana yang dimulai dengan huruf atau garis bawah dan hanya berisi huruf, angka, dan garis bawah. Sink tidak mendukung identifier yang diapit tanda kutip atau karakter khusus, seperti tanda hubung, dalam nama tabel atau skema.
Pengoptimalan Performa
Batching dan backpressure
Flush dipicu ketika salah satu kondisi terpenuhi:
- Buffer mencapai
batchsizebaris, yang nilai defaultnya adalah1000. - Umur buffer melebihi
batchinterval, yang nilai defaultnya adalah100 milliseconds.
Ketika basis data tidak dapat mengimbangi laju data yang masuk, sink meneruskan backpressure ke hulu menuju sumber.
Panduan latensi dan throughput:
- Untuk beban kerja latensi rendah dengan mode real time, kurangi
batchintervaluntuk menjamin waktu maksimum yang lebih singkat sebelum pembilasan. Lihat konsep mode waktu nyata untuk konsep dan contoh mode waktu nyata untuk contoh kode. - Untuk beban kerja dengan throughput tinggi, tingkatkan
batchsizeguna mengurangi overhead pada setiap transaksi.
Perilaku koneksi
Sink menggunakan pengumpulan koneksi pada pelaksana. Secara default, setiap tugas menggunakan satu koneksi database.
Databricks merekomendasikan agar Anda menggunakan nilai default 1 task untuk setiap koneksi. Jika Anda meningkatkan jumlah tugas untuk setiap koneksi, Anda dapat menyebabkan ketidakcocokan koneksi dan meningkatkan latensi untuk koneksi throughput tinggi.
Untuk mengonfigurasi rasio tugas terhadap koneksi, atur konfigurasi Spark spark.databricks.sql.streaming.jdbc.tasksPerConnection. Jika database target memiliki batas koneksi rendah, kurangi jumlah partisi acak atau tingkatkan spark.databricks.sql.streaming.jdbc.tasksPerConnection.
Sink secara otomatis mencoba kembali kesalahan JDBC sementara, termasuk kegagalan koneksi, kebuntuan, dan pembatasan laju. Jika sink menghabiskan semua upaya percobaan ulang, kueri akan gagal.
Pemicu dan mode output yang didukung
Triggers
Tabel ini memperlihatkan dukungan untuk jenis pemicu Streaming Terstruktur:
| Pemicu | Dukungan |
|---|---|
realTime |
Yes |
ProcessingTime |
Yes |
AvailableNow |
Yes |
Once |
Yes |
Mode keluaran
Tabel ini memperlihatkan dukungan untuk mode output Streaming Terstruktur:
| Mode keluaran | Dukungan |
|---|---|
update |
Yes |
append |
Yes. Perilaku identik dengan update. Kueri melakukan upsert saat tabel target memiliki kunci primer; jika tidak, kueri melakukan penyisipan. Lihat Perilaku upsert. |
complete |
No |
Batasan
- Komputasi tanpa server dan alur Lakeflow tidak didukung.
- Hanya Lakebase yang didukung sebagai target tulis. Database eksternal yang kompatibel dengan PostgreSQL tidak didukung.