關於使用 Kafka 搭配 Azure Databricks 的常見問題。
為什麼我會收到 Kafka 選項不支援或無法被辨識的錯誤訊息?
此錯誤發生在設定 Kafka 用戶端設定選項時忘記使用前綴。kafka. 所有直接傳送給 Kafka 用戶端的選項必須在前加上 kafka.:
以下程式碼顯示缺少前綴的錯誤選項 kafka. :
.option("security.protocol", "SASL_SSL")
.option("sasl.mechanism", "PLAIN")
以下代碼顯示正確的選項:
.option("kafka.security.protocol", "SASL_SSL")
.option("kafka.sasl.mechanism", "PLAIN")
Spark Kafka 連接器的選項(例如 subscribe、 startingOffsets、 maxOffsetsPerTrigger)不需要前綴。 完整選項列表請參見 卡夫卡。
為什麼我會收到有關封裝處理後卡夫卡類別的錯誤?
Azure Databricks需要使用著色的 Kafka 類別(前綴為 kafkashaded. 或 shadedmskiam.)。 如果您看到錯誤,例如 RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED,則必須使用反白的類別名稱:
-
org.apache.kafka.*類別需要kafkashaded.前綴。 例如:kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule -
software.amazon.msk.*類別需要shadedmskiam.前綴。 例如:shadedmskiam.software.amazon.msk.auth.iam.IAMLoginModule
為什麼我在連接卡夫卡時會收到一個錯誤或代碼?TimeoutException
常見的原因包括:
- 網路連線:運算叢集無法連上 Kafka 代理伺服器。 檢查防火牆規則、安全群組和 VPC 設定。
-
錯誤的啟動伺服器:請確認
kafka.bootstrap.servers主機名稱和埠口是否正確。 - DNS 解析:確認 Kafka 經紀人主機名稱能否從 Azure Databricks 網路解析。
- SSL/TLS 問題:若使用 SSL,請確認憑證是否正確設定。
針對 Private Link 或 VPC 對等互連設定,請確認已設定正確的網路路由。
我應該用批次模式還是串流模式來玩 Kafka?
這取決於你的使用情境:
-
串流模式 (
spark.readStream):用於需要持續資料處理或低延遲擷取時。 -
批次模式 (
spark.read):用於一次性資料載入、回填或除錯。 需要同時startingOffsets和endingOffsets。
請參閱 「配置結構化串流觸發間隔」 ,了解如何設定觸發間隔,如 AvailableNow、 ProcessingTime、 及 即時模式。
我可以在一個串流中閱讀多個卡夫卡主題嗎?
是的,你可以使用:
-
subscribe:提供一個逗號分隔的主題列表,例如.option("subscribe", "topic1,topic2")。 -
subscribePattern:使用Java正則表達式模式來匹配主題名稱,例如.option("subscribePattern", "topic-.*")。
如何在 Lakeflow 管線中使用 Kafka?
Lakeflow 管線內建對 Kafka 來源的支援。
你可以定義一個串流表,從 Kafka 讀取,如以下程式碼所示:
Python
import dlt
@dlt.table
def kafka_bronze():
return (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.load()
)
SQL
CREATE OR REFRESH STREAMING TABLE kafka_bronze AS
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:port>',
subscribe => '<topic>'
);
欲了解更多關於 Lakeflow 管線中串流來源的詳細資訊,請參閱 「管線中的負載資料 」。
我該如何反序列化 Kafka 的鍵和值欄位?
key 與 value 欄位會以 BINARY 型別傳回。 使用 DataFrame 操作根據你的資料格式反序列化:
-
字串資料:用
cast("string")來將二進位轉換成字串。 -
JSON 資料:將資料轉型為字串後,使用
from_json()。 請參閱from_json函式。 -
Avro 資料:用於
from_avro()反序列化 Avro 編碼的資料。 請參閱 讀取和寫入串流 Avro 數據。 -
協定緩衝區:用於
from_protobuf()反序列化 protobuf 資料。 請參閱 讀寫協議緩衝區。
為什麼我會出現冪等寫錯誤?
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_DATA_LOSS_ERROR,我該怎麼解決它呢?
當 Kafka 來源偵測到檢查點中儲存的偏移量不再可用時,通常會發生這個錯誤,原因通常如下:
- 資料流暫停的時間超過了 Kafka 保留期間。
- 卡夫卡主題資料被刪除或重新建立。
- Kafka 經紀人遭遇資料遺失。
解決方法:
-
若資料遺失可接受:設定
.option("failOnDataLoss", "false")以允許串流從最早可用位置繼續。 -
若資料遺失不可接受:重置檢查點並從偏移量重新處理
earliest,或還原遺失的 Kafka 資料。
詳情請參見 KAFKA_DATA_LOSS錯誤條件 。
我要如何控制從 Kafka 讀取資料的速度?
可使用 maxOffsetsPerTrigger 選項限制每個微批次處理的偏移量(約等於記錄數量)。 這有助於避免大量批次導致下游處理過載或在補上積壓時造成記憶體問題。
Python
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.option("maxOffsetsPerTrigger", 10000)
.load()
)
Scala
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.option("maxOffsetsPerTrigger", 10000)
.load()
SQL
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:port>',
subscribe => '<topic>',
maxOffsetsPerTrigger => '10000'
);
或者,也可以使用像 minPartitions 或 maxRecordsPerPartition 這樣的選項來控制每批次建立多少 Spark 分割區。
我該如何監控我的串流與最新的 Kafka 偏移相差多少?
使用在串流查詢進度中可用的 avgOffsetsBehindLatest、maxOffsetsBehindLatest 和 minOffsetsBehindLatest 指標。 這些資料會顯示你的串流在所有訂閱主題分割中落後最新可用偏移量的多少。 請參見在 Azure Databricks 上監控結構化串流查詢。
你也可以用 estimatedTotalBytesBehindLatest 來估算尚未處理的資料總位元組數。
為什麼我的 Kafka 偏移延遲指標在升級到 Databricks Runtime 17.1 後,會顯示持續的非零值?
在 Databricks Runtime 17.1 及以上版本中,每個微批次完成後會擷取最新的 Kafka 偏移量。 在持續接收資料的主題中,積壓指標可能顯示小且持續的非零值。 這是預期中的行為,並不代表溪流正在落後。
在 Databricks Runtime 17.0 及以下版本中,最新的 Kafka 偏移量會在微批次開始時擷取。 當串流查詢持續消耗微批次開始時所有可用紀錄時,積壓指標可能會回傳 0 。
如果數值很大或持續成長,資料流可能無法跟上輸入資料。 請參見在 Azure Databricks 上監控結構化串流查詢。
為什麼我的 Kafka 串流初始化很慢?
Kafka Streams 的執行需要一些時間:
- 連接到 Kafka 叢集並取得元資料。
- 探索主題分割區。
- 取得初始偏移量。
對於本地或遠端的 Kafka 叢集,網路延遲會顯著影響初始化時間。 如果你正在運行觸發/排程的管線,且經常重啟,建議使用連續串流模式以避免重複初始化的負擔。
為什麼增加更多 Spark 執行器不會增加我的 Kafka 吞吐量?
一旦 Kafka 經紀人飽和,增加更多 Spark 執行器會增加成本卻不增加吞吐量。
卡夫卡成為瓶頸的徵兆:
- 儘管增加了更多核心,吞吐量仍然達到瓶頸。
- Kafka 經紀商的 CPU 或網路利用率很高。
- Spark 任務快速完成,但請等待新資料。
為了解決這個問題,可以透過增加 brokers 或增加分割區數量來擴展 Kafka 叢集,以分散負載。
我該如何優化 Kafka 串流的成本與運算利用率?
針對微批次與 AvailableNow 模式:
- 調整叢集規模:監控指標並設定適合尖峰負載的固定叢集大小。
-
使用
maxOffsetsPerTrigger:限制批次大小以控制負載高峰時的資源使用。 - 避免自動擴展:串流工作持續執行,新增或移除節點會造成任務重新平衡的開銷。
-
減少資料偏斜:偏斜的分割區會讓某些任務處理顯著較多的資料,導致緩慢的工作,拖慢整體批次處理的完成速度,並使得閒置的任務浪費運算資源。 利用
minPartitions將大型 Kafka 分割區分割成較小的 Spark 分割區的選項,以提升處理平衡。
對於即時模式,計算規模調整尤其重要,因為任務在等待資料時可能會閒置。 主要考慮因素:
- 設定
maxPartitions使每個任務處理多個 Kafka 分區,以降低開銷。 - 調整
spark.sql.shuffle.partitions適用於洗牌密集的工作。
關於即時模式叢集大小的指引,請參見 「計算尺寸 」。
為何我的資料流沒有回傳任何資料,儘管在主題中存在資料?
常見的原因包括:
-
設定錯誤
startingOffsets:預設值為latest,只讀取串流開始後的新資料。 設定startingOffsets為earliest讀取現有資料。 - 主題名稱錯誤:請確認你訂閱的是正確的主題。
- 認證問題:你的Stream可能成功連線,但缺乏讀取主題的權限。 檢查你的 Kafka ACL。
-
偏移過期:如果你的直播被停止很久,且檢查點的偏移已過期(被卡夫卡保留刪除),你可能需要重置檢查點或調整
failOnDataLoss。