Berlangganan Google Pub/Sub

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:

  • clientEmail
  • clientId
  • privateKey
  • privateKeyId

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.