Hubungkan ke Lakebase

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. catalog adalah 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, misalnya id atau user_id,event_type. Jika Anda menghilangkan upsertkey, 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 format projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Lihat Pengidentifikasi komputasi.
  • <database>: Opsional. Nama database Postgres yang menjadi target. Secara default menjadi databricks_postgres. Lihat Mengelola database.
  • <schema>.<table>: Tabel target dalam format schema.table. Jika Anda menghilangkan skema, sink menggunakan skema public. 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, misalnya id atau user_id,event_type. Jika Anda menghilangkan upsertkey, 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 KEY dengan 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 batchsize baris, yang nilai defaultnya adalah 1000.
  • Umur buffer melebihi batchinterval, yang nilai defaultnya adalah 100 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 batchinterval untuk 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 batchsize guna 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.