這很重要
這項功能目前處於 公開預覽版。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。
串流代表外部串流資料來源,例如 Apache Kafka。 串流會儲存連線詳細資料、驗證、結構描述及資料擷取設定。 串流建立後,你可以使用 功能檢視 定義來參考它,建立即時串流功能。
溪流有三部分名稱(catalog.schema.stream_name)。 串流的存取權限是由其關聯的擷取資料表所決定。 詳情請參見 「攝取與回填 」。
要求
- 用於執行筆記本指令:無伺服器或運行 Databricks Runtime 17.0 ML 或以上版本的經典運算叢集。
-
feature-engineering-client必須安裝 Python 套件 0.16.0 或以上版本。
建立數據流
用來 create_stream() 建立新的串流。 一個串流需要四個配置元件:
- 來源設定:指定串流平台(例如 Kafka)及來源特定細節(例如 Kafka 的主題訂閱)。
- 連線設定:規定如何連接及驗證串流平台,包括啟動伺服器與憑證。
- Schema config:定義訊息鍵與值的結構。
- 擷取設定:指定串流資料的擷取地點與方式。 詳情請參見 「攝取與回填 」。
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
StreamBackfillSource,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="events-topic"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "transaction_id": {"type": "string"},'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string", "format": "date-time"}'
' }'
'}'
)
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
),
)
連接串流來源
在定義串流功能之前,請先連線至你的 Kafka broker,並測試串流 Lakeflow 管線連線。 請參考 在無伺服器運算上進行串流及 連線至 Apache Kafka。
關於 AWS 管理串流(Amazon MSK),請參見 Amazon MSK 的無伺服器私密連線。 關於 Kafka 認證選項的詳細資訊,請參見認證。
Authentication
Unity Catalog 連線(建議)
使用 Unity Catalog 連線來認證你的 Kafka 叢集。 這是管理式認證的推薦方法。 要建立連結,請參見 「建立連結」。
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
直接 mTLS
若要使用直接 mTLS 驗證,請提供儲存在 Unity Catalog 磁碟區中的金鑰儲存庫和信任儲存庫檔案,其密碼透過 Databricks secret scope 參照。 欲了解更多關於使用 Kafka 的 SSL 認證資訊,請參閱「使用 SSL 連接 Azure Databricks 與 Kafka」。
from databricks.feature_engineering.entities import (
DirectMtlsConfig,
MtlsConfig,
SecretScopeReference,
)
connection_config = DirectMtlsConfig(
bootstrap_servers="broker1:9092,broker2:9092",
mtls_config=MtlsConfig(
keystore_location="/Volumes/my_catalog/my_schema/my_volume/keystore.jks",
keystore_password_ref=SecretScopeReference(
scope="my_scope", key="keystore_password"
),
key_password_ref=SecretScopeReference(
scope="my_scope", key="key_password"
),
truststore_location="/Volumes/my_catalog/my_schema/my_volume/truststore.jks",
truststore_password_ref=SecretScopeReference(
scope="my_scope", key="truststore_password"
),
),
)
SASL
預覽版期間不支援 SASL 認證(SASL/SCRAM 與 SASL/PLAIN)。
訂閱模式
訂閱模式指定 Stream 如何選擇要取用的 Kafka 主題。 支援三種模式:
| Mode | Description | 範例 |
|---|---|---|
subscribe |
逗號分隔主題名稱列表 | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
以 Java 正規表示式比對主題名稱 | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
指定主題與分割區指派的 JSON | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
模式組態
使用 JSON Schema 格式定義訊息鍵與值的結構。 對於卡夫卡來源, payload_schema 對應卡夫卡訊息值( value 卡夫卡鍵值模型中的值值),也 key_schema 對應卡夫卡訊息鍵。 至少必須提供其中一項 payload_schema 或 key_schema 。
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"}'
' }'
'}'
)
),
key_schema=SchemaConfig(
json_schema='{"type": "string"}'
),
)
若未提供金鑰或有效載荷的結構,則視為簡單字串。
攝取與回填
ingestion_config 參數用於設定如何擷取及儲存串流資料,以供訓練和推論服務之用。
對串流的存取由擷取資料表控管:
-
SELECT在擷取資料表上授與對串流的讀取權限。 -
MANAGE會對擷取資料表授與刪除權限。
欲了解更多資料表權限資訊,請參閱 資料表 與 Unity 目錄權限參考資料。
資料輸入管線
當串流被建立時,Databricks 會啟動一個受管理的擷取管線,持續讀取 Kafka 主題的訊息並將其寫入 Delta 表(即擷取表)。 管線會從最新的 Kafka offset 開始,並持續執行,僅擷取在串流建立後傳入的新訊息。 此擷取表用於 訓練串流功能。 當串流被刪除時,其擷取管線和擷取表也會被刪除。
攝取地點
ingestion_destination 指定串流資料寫入的三部分組成 Delta 資料表名稱。
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
攝取表結構
擷取表包含訊息資料及元資料欄位:
| Column | 類型 | Description |
|---|---|---|
key |
視情況而定(從 key_schema 開始) |
卡夫卡訊息金鑰,依照你提供的結構設計。 |
value |
視情況而定(從 payload_schema 開始) |
Kafka 訊息值(有效載荷),其結構依照你提供的結構描述。 |
stream_record_timestamp |
TIMESTAMP |
紀錄時間戳記。 對於前向填補資料,這是 Kafka broker 的擷取時間戳記。 回填資料則由客戶提供。 |
kafka_topic |
STRING |
這張唱片就是從卡夫卡主題中取材的。 |
kafka_partition |
INT |
記錄是從卡夫卡分割區被消費的。 |
kafka_offset |
LONG |
記錄在其分割區內的卡夫卡偏移量。 |
record_source |
STRING |
要麼是 "stream"(從即時 Kafka 串流向前填補),要麼是 "backfill"(從回填來源)。 |
回填源
由於前向填補管線是從最新的 Kafka offset 開始,因此無法擷取在串流建立之前就已存在的訊息。 為了讓訓練涵蓋歷史資料,請設定選用的回填來源。
當設定回填來源時,Databricks 會執行一次性的MERGE INTO作業,將回填資料列以 record_source="backfill" 複製到擷取資料表中。 合併程序僅在重疊檢查器確認回填來源與前填流的時間戳重疊後執行(參見 回填與直播資料重疊)。 若兩天內未達成重疊條件,合併程序仍會執行以避免無限期阻塞。
回填表必須包含 stream_record_timestamp 一欄,該欄位 TIMESTAMP 屬於 UTC 時區。 其他 Kafka 中繼資料欄位(kafka_topic、kafka_partition、kafka_offset)如果存在於回填來源中,則會一併傳遞;否則會設為 NULL。
from databricks.feature_engineering.entities import StreamBackfillSource
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
backfill_source=StreamBackfillSource(
delta_table_name="my_catalog.my_schema.historical_events"
),
)
回填資料與即時串流資料之間的重疊
在回填與擷取資料表之間執行 MERGE 之前,會先進行重疊檢查,比較這兩個資料表中的時間戳:
-
回填最大值:回填來源的最大值
stream_record_timestamp。 -
擷取最小值:擷取資料表中資料列的最小
stream_record_timestamp(record_source="stream")。
當回填物的最新時間戳記比攝取表最早的時間戳至少多 1 小時時,合併程序即會繼續。 這種重疊確保攝取表中不會有空隙。 若兩天內未達成重疊條件,合併程序仍會執行以避免無限期阻塞。
由於資料擷取管線是從最新的 Kafka offset 開始,因此只會擷取在串流建立後才到達的訊息。 你的回填來源必須包含延伸至攝取時間範圍內的資料,而不只是延伸到串流建立時間。
例如,如果你在下午 3:00 建立串流,向前填補管線會從下午 3:00 起開始讀取訊息。 你的回填來源必須包含時間戳記至少涵蓋到下午 4:00 的資料(即向前填補開始後 1 小時),才能通過重疊檢查。 這表示你應該在下午4點後更新回填表,確保攝取表沒有空隙。
Deduplication
用 deduplication_columns 來指定欄位路徑,以便在回填與前填資料資料的擷取過程中辨識重複列。 巢狀欄位使用點符號(例如 "value.user_id")。
根據您的資料選擇去重複欄位:
- 如果串流中的每筆記錄都包含唯一識別碼(例如),
value.transaction_id請使用該欄位進行重複資料刪除。 - 如果你的回填來源包含
kafka_partition和 欄位kafka_offset,請用它們來唯一識別每筆紀錄。 - 若未指定去重複欄位,預設的去重複索引鍵為
key、value和stream_record_timestamp的完整組合。 這不建議,因為這種嚴格的標準匹配很容易導致重複。
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
管理串流
去直播
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
列表串流
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
設定 include_schemas=True 包含完整架構細節。 結構可能很大,這可能導致操作持續時間較長。 若要單獨取得結構,請使用 get_stream。
刪除串流
刪除串流同時也會刪除其擷取管線和擷取表。
Warning
任何參考已刪除串流的模型或特徵將無法存取底層串流資料。 如果你需要這些資料但不再需要串流,請先在刪除前建立一個擷取表的副本。
client.delete_stream(name="my_catalog.my_schema.my_stream")
範例筆記本
關於建立串流、定義串流功能並部署到服務端端的端對端範例,請參考以下筆記本: