連接 Apache Kafka

本頁說明如何在 Azure Databricks 上執行結構化串流工作負載時,使用 Apache Kafka 作為來源或匯入工具。

欲了解更多關於 Kafka 的資訊,請參閱 Apache Kafka 文件。

從 Kafka 讀取資料

使用該 kafka 格式來設定與 Kafka 的連線。 以下是串流閱讀的範例:

Python

df = (spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "latest")
  .load()
)

Scala

val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "latest")
  .load()

SQL

CREATE OR REFRESH STREAMING TABLE <table_name> AS
SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<server:ip>',
  subscribe => '<topic>'
);

Azure Databricks 也支援 Kafka 的批次讀取,如下範例所示:

Python

df = (spark.read
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load()
)

Scala

val df = spark.read
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load()

SQL

SELECT * FROM read_kafka(
  bootstrapServers => '<server:ip>',
  subscribe => '<topic>',
  startingOffsets => 'earliest',
  endingOffsets => 'latest'
);

對於累加式批次載入,Databricks 建議將 Kafka 與 Trigger.AvailableNow 搭配使用。 參見 AvailableNow:增量批次處理。

在 Databricks Runtime 13.3 LTS 及以上版本中,Azure Databricks 也提供一個 SQL 函式來讀取 Kafka 資料。 只有在 Lakeflow 管線中,或在 Databricks SQL 中搭配串流資料表時,才支援使用 SQL 進行串流處理。 請參閱 read_kafka 資料表值函式。

設定 Kafka 結構化串流讀取器

無論是批次查詢還是串流查詢,你都必須設定 Kafka 原始碼的啟動伺服器,並有以下選項:

鑰匙 價值 說明
kafka.bootstrap.servers 一個逗號分隔的 host:port 列表 Kafka 叢集引導伺服器

要設定訂閱主題,您必須指定以下選項之一:

Option 價值 說明
subscribe 以逗號分隔的主題清單。 要訂閱的主題清單。
subscribePattern Java 正則表達式字串。 用於訂閱主題的範式。
assign JSON 字串 {"topicA":[0,1],"topic":[2,4]}。 針對 topicPartitions 特定消費量。

完整選項請參見 卡夫卡 。

Kafka 列的架構

Kafka 結構化串流閱讀器回傳的列包含以下結構:

資料行 類型
key binary
value binary
topic string
partition int
offset long
timestamp timestamp
timestampType int

key和value 一律以ByteArrayDeserializer 還原序列化為位元組陣列。 使用 DataFrame 操作(例如 cast("string") 或 from_avro)來明確反序列化鍵與值。

將資料寫入 Kafka

以下是串流寫入 Kafka 的範例:

Python

(df.writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .start()
)

Scala

df.writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .start()

Azure Databricks 也支援批次寫入 Kafka 資料匯的語意,如下範例所示:

Python

(df.write
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .save()
)

Scala

df.write
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .save()

設定 Kafka 結構化串流寫入器

這很重要

Databricks Runtime 13.3 LTS 和更新版本包含了一個更新的 kafka-clients 程式庫,該程式庫預設啟用等冪寫入功能。 如果 Kafka 接收器使用 2.8.0 或更低版本,並且已設定 ACL 但未啟用 IDEMPOTENT_WRITE,則寫入將會失敗並出現錯誤訊息 org.apache.kafka.common.KafkaException:Cannot execute transactional method because we are in an error state。

透過升級至 Kafka 2.8.0 或更新版本或者在設定結構化串流寫入器時設定 .option(“kafka.enable.idempotence”, “false”),來解決此錯誤。

以下是常見的寫信給 Kafka 的選項:

鑰匙 價值 預設值 說明
kafka.boostrap.servers 以逗號分隔的 <host:port> 清單 沒有 Required. Kafka bootstrap.servers 組態。
topic STRING 未設定 Optional. 設定所有列要寫的主題。 此選項會覆寫資料中已存在的任何主題欄位。
includeHeaders BOOLEAN false Optional. 是否要在數據列中包含 Kafka 標頭。

完整選項請參見 卡夫卡水槽 。

卡夫卡作家的架構

在寫入 Kafka 資料時,所提供的 DataFrame 可能包含以下欄位:

欄位名稱 必要或選擇性 類型
key optional STRING 或 BINARY
value required STRING 或 BINARY
headers optional ARRAY
topic 選擇性(如果 topic 設為 writer 選項,則會忽略) STRING
partition optional INT

驗證

Azure Databricks 支援多種 Kafka 認證方法,包括 Unity Catalog 服務憑證、SASL/SSL,以及針對 AWS MSK、Azure 事件中樞 和 Google Cloud Managed Kafka 的雲端專用選項。 請參閱驗證。

擷取 Kafka 指標

要監控串流查詢落後 Kafka 的延遲,請使用 avgOffsetsBehindLatest、 maxOffsetsBehindLatest和 minOffsetsBehindLatest 指標。 這些指標顯示所有已訂閱主題分區相對於 Kafka 中最新偏移量的平均、最大及最小偏移落後。 請參閱以互動方式讀取計量。

備註

在 Databricks Runtime 17.1 及以上版本中,每個微批次完成後會擷取最新的 Kafka 偏移量。 在持續接收資料的主題中,積壓指標可能顯示小且持續的非零值。 這是預期中的行為,並不代表溪流正在落後。

在 Databricks Runtime 17.0 及以下版本中,最新的 Kafka 偏移量會在微批次開始時擷取。 當串流查詢持續消耗微批次開始時所有可用紀錄時,積壓指標可能會回傳 0 。

要估算查詢剩餘資料,請使用 estimatedTotalBytesBehindLatest 指標。 此指標根據過去 300 秒內處理的批次,估計所有訂閱分割區剩餘的位元組總數。 你可以透過設定選項 bytesEstimateWindowLength 來修改這個估算的時間範圍。

例如,要將視窗長度設定為 10 分鐘:

Python

df = (spark.readStream
  .format("kafka")
  .option("bytesEstimateWindowLength", "10m") # m for minutes, you can also use "600s" for 600 seconds
)

Scala

val df = spark.readStream
  .format("kafka")
  .option("bytesEstimateWindowLength", "10m") // m for minutes, you can also use "600s" for 600 seconds

如果在筆記本中執行串流,您可以在串流查詢進度儀表板中的 [原始資料] 索引標籤下查看這些計量:

{
  "sources": [
    {
      "description": "KafkaV2[Subscribe[topic]]",
      "metrics": {
        "avgOffsetsBehindLatest": "4.0",
        "maxOffsetsBehindLatest": "4",
        "minOffsetsBehindLatest": "4",
        "estimatedTotalBytesBehindLatest": "80.0"
      }
    }
  ]
}

更多資訊,請參閱 在 Azure Databricks 上監控結構化串流查詢。

Kafka 到 Delta Lake 的範例

以下範例展示了使用 觸發 availableNow 器從 Kafka 到三角湖資料表的增量串流寫入的完整工作流程。 你可以用這種方法來處理增量式資料擷取的工作負載。

這個範例使用固定的 JSON 架構。 對於像 Avro 或 Protobuf 這類格式,可以使用 from_avro 或 from_protobuf。 你也可以整合 schema 登錄檔。 請參考 結構登錄器的範例。

Python

from pyspark.sql.functions import from_json, col

# Define simple JSON schemas for key and value
key_schema = "user_id STRING"
value_schema = "event_type STRING, event_ts TIMESTAMP"

# Configure Kafka options with service credentials
kafka_options = {
  "kafka.bootstrap.servers": "<bootstrap-server>:9092",
  "subscribe": "<topic-name>",
  "databricks.serviceCredential": "<service-credential-name>",
}

# Read from Kafka and parse JSON
parsed_df = (spark.readStream
  .format("kafka")
  .options(**kafka_options)
  .load()
  .select(
    from_json(col("key").cast("string"), key_schema).alias("key"),
    from_json(col("value").cast("string"), value_schema).alias("value")
  )
  .select("key.*", "value.*")
)

# Write to Delta table
query = (parsed_df.writeStream
  .format("delta")
  .option("checkpointLocation", "/path/to/checkpoint")
  .trigger(availableNow=True)
  .toTable("catalog.schema.events_table")
)

query.awaitTermination()

Scala

import org.apache.spark.sql.functions.{from_json, col}
import org.apache.spark.sql.streaming.Trigger
import org.apache.spark.sql.types.StructType

// Define JSON schemas for key and value
val keySchema = "user_id STRING"
val valueSchema = "event_type STRING, event_ts TIMESTAMP"

// Configure Kafka options with service credentials
val kafkaOptions = Map(
  "kafka.bootstrap.servers" -> "<bootstrap-server>:9092",
  "subscribe" -> "<topic-name>",
  "databricks.serviceCredential" -> "<service-credential-name>"
)

// Read from Kafka and parse JSON
val parsedDF = spark.readStream
  .format("kafka")
  .options(kafkaOptions)
  .load()
  .select(
    from_json(col("key").cast("string"), StructType.fromDDL(keySchema)).alias("key"),
    from_json(col("value").cast("string"), StructType.fromDDL(valueSchema)).alias("value")
  )
  .select("key.*", "value.*")

// Write to Delta table
val query = parsedDF.writeStream
  .format("delta")
  .option("checkpointLocation", "/path/to/checkpoint")
  .trigger(Trigger.ProcessingTime("10 seconds"))
  .toTable("catalog.schema.events_table")

query.awaitTermination()

SQL

-- Create a streaming table from Kafka using read_kafka
CREATE OR REFRESH STREAMING TABLE catalog.schema.events_table AS
SELECT
  key::string:user_id AS user_id,
  value::string:event_type AS event_type,
  to_timestamp(value::string:event_ts) AS event_ts
FROM STREAM read_kafka(
  bootstrapServers => '<bootstrap-server>:9092',
  subscribe => '<topic-name>',
  serviceCredential => '<service-credential-name>'
);

備註

在 Databricks 的無伺服器運算中,建議使用 availableNow 這個觸發器來進行增量串流。 對於低延遲的連續串流,請使用 Lakeflow pipelines 連續模式。 請參閱 結構化串流觸發器 以了解完整的支援選項清單。