本頁說明如何在 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 連續模式。 請參閱 結構化串流觸發器 以了解完整的支援選項清單。