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.
Gunakan konektor bawaan untuk berlangganan Google Pub/Sub. Konektor ini memiliki semantik pemrosesan tepat satu kali untuk baris data dari subscriber.
Catatan
Pub/Sub mungkin mengirimkan baris duplikat, atau baris dapat sampai ke subscriber dalam urutan yang tidak berurutan. Anda harus menulis kode untuk menangani baris duplikat dan tidak berurutan.
Mengonfigurasi aliran Pub/Sub
Contoh kode berikut menunjukkan cara mengonfigurasi Bacaan Streaming Terstruktur dari Pub/Sub dan mengautentikasi dengan kunci privat.
Python
auth_options = {
"clientId": client_id,
"clientEmail": client_email,
"privateKey": private_key,
"privateKeyId": private_key_id
}
query = (spark.readStream
.format("pubsub")
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.options(auth_options)
.load()
)
Scala
val authOptions: Map[String, String] =
Map("clientId" -> clientId,
"clientEmail" -> clientEmail,
"privateKey" -> privateKey,
"privateKeyId" -> privateKeyId)
val query = spark.readStream
.format("pubsub")
// Creates a Pub/Sub subscription if one does not already exist with this ID
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.options(authOptions)
.load()
SQL
CREATE OR REFRESH STREAMING TABLE pubsub_raw
AS SELECT * FROM STREAM read_pubsub(
subscriptionId => 'mysub',
projectId => 'myproject',
topicId => 'mytopic',
clientEmail => secret('pubsub-scope', 'clientEmail'),
clientId => secret('pubsub-scope', 'clientId'),
privateKeyId => secret('pubsub-scope', 'privateKeyId'),
privateKey => secret('pubsub-scope', 'privateKey')
);
Untuk opsi konfigurasi lainnya, lihat Mengonfigurasi opsi untuk pembacaan streaming Pub/Sub.
Mengonfigurasi akses ke Pub/Sub
Kredensial Anda harus memiliki peran berikut:
| Peran | Diperlukan atau opsional | Bagaimana peran digunakan |
|---|---|---|
roles/pubsub.viewer atau roles/viewer |
Wajib | Memeriksa apakah langganan ada dan mendapatkan langganan. |
roles/pubsub.subscriber |
Wajib | Mengambil data dari langganan. |
roles/pubsub.editor atau roles/editor |
Opsional | Mengaktifkan pembuatan langganan jika tidak ada dan memungkinkan penggunaan deleteSubscriptionOnStreamStop untuk menghapus langganan pada penghentian streaming. |
Catatan
Jika Anda memberikan roles/pubsub.viewer dan roles/pubsub.subscriber di tingkat sumber daya daripada tingkat proyek, Anda harus menerapkan kedua peran ke topik dan langganan. Jika Anda tidak menggunakan peran roles/pubsub.editor atau roles/editor yang opsional, memberikan peran yang diperlukan hanya pada topik saja tidaklah cukup.
Databricks merekomendasikan agar Anda menggunakan rahasia saat menggunakan kunci. Opsi berikut diperlukan untuk mengotorisasi koneksi:
clientEmailclientIdprivateKeyprivateKeyId
Memahami skema Pub/Sub
Skema untuk stream sesuai dengan baris yang diambil dari Pub/Sub, seperti yang dijelaskan dalam tabel berikut:
| Bidang | Jenis |
|---|---|
messageId |
StringType |
payload |
ArrayType[ByteType] |
attributes |
StringType |
publishTimestampInMillis |
LongType |
Mengonfigurasi opsi untuk baca streaming Pub/Sub
Beberapa opsi konfigurasi Pub/Sub menggunakan konsep pengambilan data alih-alih mikro-batch. Ini merupakan detail implementasi internal, dan opsinya bekerja serupa dengan konektor Structured Streaming lainnya, kecuali bahwa baris data diambil terlebih dahulu lalu diproses.
Untuk daftar lengkap opsi, lihat Pub/Sub.
Menggunakan pemrosesan batch inkremental dengan Pub/Sub
Anda dapat menggunakan Trigger.AvailableNow untuk memproses baris yang tersedia dari sumber Pub/Sub sebagai batch inkremental.
Azure Databricks mencatat tanda waktu ketika Anda memulai pembacaan dengan pengaturan Trigger.AvailableNow. Baris data yang diproses oleh batch mencakup semua data yang sebelumnya telah diambil serta semua baris data yang baru dipublikasikan dengan stempel waktu yang lebih kecil daripada stempel waktu mulai yang tercatat. Untuk informasi selengkapnya, lihat AvailableNow: Pemrosesan batch inkremental.
Memantau metrik streaming Pub/Sub
Metrik kemajuan Streaming Terstruktur melaporkan jumlah baris yang diambil dan siap diproses, ukuran baris yang diambil dan siap diproses, dan jumlah duplikat yang terlihat sejak streaming dimulai.
Berikut ini adalah contoh metrik Pub/Sub:
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}
Batasan
Pub/Sub tidak mendukung eksekusi spekulatif dengan spark.speculation.