設定串流

這很重要

這項功能目前處於 公開預覽版。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。

串流代表外部串流資料來源,例如 Apache Kafka。 串流會儲存連線詳細資料、驗證、結構描述及資料擷取設定。 串流建立後,你可以使用 功能檢視 定義來參考它,建立即時串流功能。

溪流有三部分名稱(catalog.schema.stream_name)。 串流的存取權限是由其關聯的擷取資料表所決定。 詳情請參見 「攝取與回填 」。

要求

  • 用於執行筆記本指令:無伺服器或運行 Databricks Runtime 17.0 ML 或以上版本的經典運算叢集。
  • 必須安裝 0.18.0 或以上版本的 feature-engineering-client Python 套件。

連接串流來源

在定義串流功能之前,請先連線至你的 Kafka broker,並測試串流 Lakeflow 管線連線。 Feature Store 依賴無伺服器 SDP,這表示你需要有某種機制,將你的傳統運算資源(代理程式或端點)連線至 Databricks 的無伺服器運算資源。 這可以透過像 privatelink 這類產品,或是讓你的經典運算系統從公共網際網路存取來實現。

建立數據流

用來 create_stream() 建立新的串流。 一個串流需要四個配置元件:

  • 來源設定:指定串流平台及特定來源細節,例如卡夫卡來源的主題訂閱。
  • 連線設定:規定如何連接及驗證串流平台,包括啟動伺服器與憑證。
  • Schema config:定義訊息鍵與值的結構。
  • 擷取設定:指定串流資料的擷取地點與方式。 詳情請參見 「攝取與回填 」。

關於特定 source_config 來源與連線設定,以及完整 create_stream() 範例,請參見 Apache Kafka。 綱要和擷取選項為所有來源共用。

阿帕契卡夫卡

要從 Apache Kafka 串流,請使用 KafkaStreamConfig 原始設定檔,並使用 Unity 目錄連線進行驗證。 如需 Kafka 連線相關資訊,請參閱 無伺服器運算上的串流處理 和 連線至 Apache Kafka。

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"
        ),
    ),
)

卡夫卡訂閱模式

訂閱模式指定 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]}')

卡夫卡認證

使用 Unity Catalog 連線來認證你的 Kafka 叢集。 這是管理式認證的推薦方法。 要建立連結,請參見 「建立連結」。 串流的建立者必須在該連線上具有 USE CONNECTION。 任何以串流作為來源將功能具體化的使用者,也必須在該連線上具有 USE CONNECTION。

connection_config = StreamConnectionConfig(
    uc_connection_name="my-kafka-connection"
)

該連線支援 IAM 認證(服務憑證)與 SASL 認證。

IAM(服務憑證)

使用 Unity Catalog 服務認證進行驗證,例如使用 IAM 連線至 Amazon MSK。 若要建立服務認證,請參閱 建立服務認證。 使用 credential 選項設定服務憑證名稱:

CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
    bootstrap_servers '<bootstrap_servers>',
    credential '<service_credential>'
)

除了需要在連線上具備 USE CONNECTION 之外,使用服務憑證的身分也需要在該憑證上具備 ACCESS。 將所參照的服務認證的 ACCESS 權限授予該串流的建立者,以及任何使用該串流將功能具體化的身分。 請參閱 授權使用服務憑證以訪問外部雲端服務。

SASL

SASL 認證使用使用者名稱和密碼。 設定 sasl_mechanism 為以下之一:

  • PLAIN
  • SCRAM-SHA-256
  • SCRAM-SHA-512

使用 user 和 password 選項提供憑證。 連線會安全地儲存這些憑證。

以下範例使用 SASL/SCRAM。 SASL/PLAIN 則設 sasl_mechanism 為 PLAIN。

CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
    bootstrap_servers '<bootstrap_servers>',
    sasl_mechanism 'SCRAM-SHA-512',
    user '<username>',
    password '<password>'
)

直接 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"
        ),
    ),
)

結構描述設定

定義訊息鍵與值的結構,使擷取與功能定義能讀取個別欄位。 對於卡夫卡來源, payload_schema 對應卡夫卡訊息值( value 卡夫卡鍵值模型中的值值),也 key_schema 對應卡夫卡訊息鍵。 至少必須提供其中一項 payload_schema 或 key_schema 。

每個格式 SchemaConfig 接受三種格式之一,與來源訊息序列化方式相符: json_schema、、 avro_schema或 proto_schema。 若未提供金鑰或有效載荷的結構,則視為簡單字串。

本節的程式碼範例使用與 內聯 DirectSchemas宣告的結構,其中結構以字串形式提供。 若要使用外部結構登錄檔管理結構,請參見 結構登錄表 以了解詳細資訊。

JSON 架構

提供一個 JSON 結構 字串到 json_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"}'
    ),
)

Avro 架構

將Avro 結構描述字串提供給avro_schema。 支援 Avro 邏輯類型,包括 timestamp-millis、 date和 decimal。

schema_config = DirectSchemas(
    payload_schema=SchemaConfig(
        avro_schema=(
            '{'
            '  "type": "record",'
            '  "name": "Event",'
            '  "fields": ['
            '    {"name": "user_id", "type": "string"},'
            '    {"name": "amount", "type": "double"},'
            '    {"name": "event_time",'
            '     "type": {"type": "long", "logicalType": "timestamp-millis"}}'
            '  ]'
            '}'
        )
    ),
)

Protobuf 架構

提供ProtoSchemaSpec給,並附上.protoproto_schema原始文字和有效負載訊息名稱。 從 ProtoSchemaSpec匯入 databricks.feature_engineering.entities 。

message_name 必須是完整限定訊息名稱,包含在 .proto 文字中宣告的 package(例如,com.example.Event,而不是 Event)。 支援 proto2 和 proto3 語法。

google.protobuf.Timestamp 而標量包裝器類型(StringValue如 、 Int32Value、 等)則被支援,且其匯入會自動解析。 其他知名類型,如 Duration、 Struct和 Any,則被拒絕;改以支援的標量或訊息形式編碼這些值。 fixed32 fixed64標量型別和map非字串鍵也未被支援。

from databricks.feature_engineering.entities import ProtoSchemaSpec

schema_config = DirectSchemas(
    payload_schema=SchemaConfig(
        proto_schema=ProtoSchemaSpec(
            schema_text=(
                'syntax = "proto3";\n'
                'package com.example;\n'
                'import "google/protobuf/timestamp.proto";\n'
                'message Event {\n'
                '  string user_id = 1;\n'
                '  double amount = 2;\n'
                '  google.protobuf.Timestamp event_time = 3;\n'
                '}'
            ),
            message_name="com.example.Event",
        )
    ),
)

使用結構解碼資料

Databricks 使用 Spark 的 from_json、from_avro 和 from_protobuf 函式來解碼每則訊息。 無論你是宣告結構為內線或從結構登錄檔解析,以下行為都適用:

  • 格式錯誤的記錄。 解碼使用模式 PERMISSIVE ,也就是說,與模式不符的記錄會解碼為空值,而非失敗串流。
  • Avro 工會。 多個紀錄類型的聯集會解碼為一個結構體,其中每種紀錄類型各有一個欄位,且每個欄位都以其對應的 Avro 紀錄名稱命名。
  • Protobuf 類型。 無符號整數解碼為更寬的有號型別(例如 uint32 to BIGINT 和 uint64 ), DECIMAL(20,0)enum 欄位解碼為其字串名稱,純量包裝型別(例如 StringValue 和 Int32Value)則解碼為包裹型別中可空的欄位。

架構登錄

架構登錄系統儲存串流製作者與消費者使用的版本架構,並隨著結構演變而強制執行相容規則。 當外部結構登錄檔被設定時,功能商店會從登錄檔讀取結構,並用它來解碼串流訊息。 使用結構描述註冊表時,你不會直接在 Stream 上宣告該結構描述。

結構登錄支援有以下限制:

  • 僅支援 Kafka 串流。
  • 僅支援 Confluent Schema 登錄檔
  • 僅支援 Avro 和 Protobuf 格式。 要讀取 JSON 訊息,請直接宣告 schema inline。 請參見 JSON 架構。
  • 每個串流只連接到一個 Confluent 主體 作為訊息值,另一個主體用於訊息鍵(如提供)。 包含多個結構記錄的串流主題不支援配置。 如果你的串流連接到包含多個結構的主題,與該主題結構不符的紀錄會 被解碼為空。

連接到結構描述登錄檔

在 Kafka Unity 目錄連線中提供登錄連結細節作為選項,並將登錄檔 API 秘密存放在 Databricks 的秘密作用域中。 串流的 run-as(執行身分)必須具有祕密範圍的 READ 權限,因為擷取管線會在執行階段讀取該祕密。 關於如何建立和設定連線,請參見 建立連線。

將 schema_registry_url、 schema_registry_api_key、 和 schema_registry_api_secret 選項加入 用於驗證的連線。 以下範例建立一個 Kafka 連線,該連線以 Unity Catalog 服務憑證驗證代理者,並以 API 金鑰驗證登錄檔:

CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
    bootstrap_servers '<bootstrap_servers>',
    credential '<service_credential>',
    schema_registry_url 'https://<registry-host>',
    schema_registry_api_key '<registry_api_key>',
    schema_registry_api_secret secret('<scope>', '<key>')
)

將 Kafka 連線上的 選項和 Stream 上的 schema_registry_api_secret 參照都設為同一個秘密。

建立一個使用 schema registry 的串流

將 a SchemaRegistryConfig 傳遞為 schema_config. 使用 api_secret_ref 參照 Registry API 密鑰,並使用 payload_schema_locator 指定訊息值的主體與格式,或使用 key_schema_locator 指定訊息鍵的主體與格式。 至少必須提供一個定位器。

請注意這裡與 結構配置 章節中直接的結構範例的差異。 使用結構描述登錄時,你不需要在 Stream 上提供內嵌的結構描述給 schema_config。 相反地,你指定一個 SchemaRegistryConfig 來識別登錄檔中結構的結構。

from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
    KafkaStreamConfig,
    KafkaSubscriptionMode,
    StreamConnectionConfig,
    SchemaRegistryConfig,
    SchemaLocator,
    SchemaLocatorConfluentSchema,
    SchemaLocatorFormat,
    SecretScopeReference,
    IngestionConfig,
    IngestionDestination,
)

client = FeatureEngineeringClient()

stream = client.create_stream(
    name="my_catalog.my_schema.my_stream",
    source_config=KafkaStreamConfig(
        subscription_mode=KafkaSubscriptionMode(subscribe="transactions"),
    ),
    connection_config=StreamConnectionConfig(
        uc_connection_name="my-kafka-connection"
    ),
    schema_config=SchemaRegistryConfig(
        api_secret_ref=SecretScopeReference(
            scope="my_scope", key="sr_api_secret"
        ),
        payload_schema_locator=SchemaLocator(
            confluent_schema=SchemaLocatorConfluentSchema(
                subject="transactions-value"
            ),
            format=SchemaLocatorFormat.FORMAT_AVRO,
        ),
    ),
    ingestion_config=IngestionConfig(
        ingestion_destination=IngestionDestination(
            delta_table_name="my_catalog.my_schema.transactions_ingestion"
        ),
    ),
)

Confluent 主體是指定範圍,用於登錄結構的版本歷史並強制相容性。 設定 subject 為相關範圍名稱,通常由 主題名稱策略決定:

  • TopicNameStrategy(預設,主旨取自主題名稱):<topic>-value 用於值,而 <topic>-key 用於鍵。 例如,主題 transactions 的值結構使用主語 transactions-value。
  • RecordNameStrategy(主體名稱衍生自結構描述的記錄名稱,且不依賴主題名稱):完整限定記錄名稱,例如 com.example.Payment。 這是 Avro 的記錄命名空間與名稱,或是訊息的封包名稱(Protobuf)。
  • TopicRecordNameStrategy (結合主題名稱與記錄名稱): <topic>-<fully-qualified-record-name>,例如 transactions-com.example.Payment。

format 是必要的。 將它設為 SchemaLocatorFormat.FORMAT_AVRO 或 SchemaLocatorFormat.FORMAT_PROTOBUF,以符合主題序列化的方式。

圖式演化

吞噬管線會在受試者啟動時解析其當前模式。 當你在結構登錄檔註冊一個新的向下相容模式版本時,執行中的管線仍會繼續使用它最初使用的版本。

對於結構登錄檔支援的串流,擷取管線會自動每隔幾小時重新啟動一次。 每次重新啟動時,它都會採用 subject 的最新結構描述版本,而新增或變更的欄位會出現在擷取資料表中。

使用直接結構描述而非結構描述登錄的串流,會使用 update_stream 來演進其結構描述。 請參閱 「更新串流」。

關於管線如何處理與目前使用架構不符的紀錄,請參見 「利用結構解碼資料」。

依類型篩選記錄

串流會用單一的鍵值結構(如提供)來解碼每筆記錄,無論你是直接指定還是使用結構登錄檔。 由於一個主題可以攜帶多種類型的紀錄,且串流可以訂閱多個主題,請用來 record_type_filter 選擇該主題中哪些記錄屬於這個串流。

提供一個 SQL 運算式,使用點記號參照已解碼的欄位,例如 value.event_type = 'transaction'。 不符合篩選條件的紀錄會被忽略。 它們不會寫入擷取資料表,也不會用於具體化。 若要為其他記錄類型建立 Stream,請另外建立一個使用不同 record_type_filter 的 Stream。

stream = client.create_stream(
    name="my_catalog.my_schema.my_stream",
    # ...source, connection, schema, and ingestion config...
    record_type_filter="value.event_type = 'transaction'",
)

即使沒有 record_type_filter,解碼也不會讓串流失效。 與配置結構不符的紀錄會被允許 解碼。 這些紀錄的解碼方式如下:

  • 放入一資料列,包含 NULL 結構預期但記錄遺漏的欄位值(JSON、Avro 和 Protobuf)。
  • 插入到包含屬於不同記錄類型之值的資料列中(僅限 Avro 和 Protobuf)。

若要判斷擷取資料表中的某列是否屬於預期的紀錄類型,請使用下列其中一項檢查:

  • 檢查某個欄位是否等於預期值,例如 value.event_type = 'transaction'(較建議用於 Avro 和 Protobuf)。
  • 檢查欄位是否不為 NULL,例如 value.activity_id IS NOT NULL。

record_type_filter當主題的紀錄類型間結構差異很大,或你想獨立管理各記錄類型的存取時,建議使用不同的串流。 為了控制成本,Databricks 建議保留少量串流,因為每個串流都有獨立的擷取管線和擷取表。 每個串流在實體化時也會使用獨立的運算資源。 你可以使用特定功能的篩選器來進行具體化。

record_type_filter 與特徵的 filter_condition不同。 record_type_filter 設定在串流中,控制哪些紀錄被擷取並開放給所有使用串流作為來源的功能;而 filter_condition 則設定在單一特徵上,並在彙總前過濾列。 欲了解更多詳情filter_condition,請參閱串流來源的篩選條件。

攝取與回填

ingestion_config 參數用於設定如何擷取及儲存串流資料,以供訓練和推論服務之用。

對串流的存取由擷取資料表控管:

  • SELECT 在擷取資料表上授與對串流的讀取權限。
  • MANAGE 會對擷取資料表授與刪除權限。

欲了解更多資料表權限資訊,請參閱 資料表 與 Unity 目錄權限參考資料。

資料輸入管線

當串流被建立時,Databricks 會啟動一個受管理的擷取管線,持續從來源串流讀取訊息並將其寫入 Delta 資料表(即擷取資料表)。 管線從來源中最新的位置開始,持續運行,只捕捉在串流建立後抵達的新訊息。 此擷取表用於 訓練串流功能。 當串流被刪除時,其擷取管線和擷取表也會被刪除。

攝取地點

ingestion_destination 指定串流資料寫入的三部分組成 Delta 資料表名稱。

ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
)

攝取表結構

擷取表包含訊息資料及元資料欄位。 每個來源都包含通用欄位;kafka_* 欄位僅存在於 Kafka 串流中。

Column 類型 Source Description
key 視情況而定(從 key_schema 開始) Common 訊息金鑰,依照你提供的結構架構來設計。
value 視情況而定(從 payload_schema 開始) Common 訊息值(payload),根據你提供的架構結構化。
stream_record_timestamp TIMESTAMP Common 紀錄時間戳記。 對於前置資料,這是來源資料的輸入時間戳記。 回填資料則由客戶提供。
record_source STRING Common 要麼是 "stream"(由即時串流前向填補),要麼是 "backfill"(來自回填來源)。
kafka_topic STRING Kafka 這張唱片就是從卡夫卡主題中取材的。
kafka_partition INT Kafka 記錄是從卡夫卡分割區被消費的。
kafka_offset LONG Kafka 記錄在其分割區內的卡夫卡偏移量。

回填源

由於前向填充管線從來源中最新的位置開始,它不會捕捉在串流建立前就存在的訊息。 為了讓訓練涵蓋歷史資料,請設定選用的回填來源。

當設定回填來源時,Databricks 會執行一次性的MERGE INTO作業,將回填資料列以 record_source="backfill" 複製到擷取資料表中。 合併程序僅在重疊檢查器確認回填來源與前填流的時間戳重疊後執行(參見 回填與直播資料重疊)。 若兩天內未達成重疊條件,合併程序仍會執行以避免無限期阻塞。

回填表必須包含 stream_record_timestamp 一欄,該欄位 TIMESTAMP 屬於 UTC 時區。 如果回填來源中有其他中繼資料欄位,則會一併傳遞;否則會設為 NULL。 對卡夫卡而言,這些分別是 kafka_topic、 kafka_partition、 kafka_offset和 。

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 小時時,合併程序即會繼續。 這種重疊確保攝取表中不會有空隙。 若兩天內未達成重疊條件,合併程序仍會執行以避免無限期阻塞。

由於資料擷取管線會從來源中的最新位置開始讀取,因此只會擷取在串流建立之後才到達的訊息。 你的回填來源必須包含延伸至攝取時間範圍內的資料,而不只是延伸到串流建立時間。

例如,如果你在下午 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"],
)

成本歸因

在 IngestionConfig 上設定 tags 和 budget_policy_id,以將該串流受管理擷取的成本歸屬於此。 Azure Databricks 會在建立 Stream 時,將其套用至擷取 Lakeflow 管線及其向前填補和回填作業。

例如,標籤限制以及如何查詢歸屬支出,請參見 「帶有標籤與無伺服器使用政策的屬性成本」。

排除串流中的欄位

用 excluded_columns 來丟棄串流中你不想被收錄的特定欄位。 排除欄位不會寫入資料擷取表,且無法被特徵引用或用於訓練。

在訊息鍵或值中用點符號指定每一欄,例如 value.user.email 或 key.account_id。 這些欄位會從解碼 key 中移除,並在 value 資料擷取、回填和實體化過程中被移除。 如果路徑指向一個結構體,其所有巢狀欄位也會被丟棄(例如, value.address 也會丟棄 value.address.city 和 value.address.zip)。

stream = client.create_stream(
    name="my_catalog.my_schema.my_stream",
    # ...source, connection, schema, and ingestion config...
    excluded_columns=["value.user.email", "value.user.ssn"],
)

使用直接結構時,排除欄位必須已存在於鍵或值架構中,否則 create_stream 會失敗。 使用結構登錄檔時,可以在欄位存在前排除它。 排除欄位也不能是重複去重欄位,因為重複欄位是識別重複列所必需的。 任何參照已排除資料欄的特徵(例如作為實體、時間序列或輸入)都會建立失敗。

你可以在建立後更改串流的排除欄位,無論是 update_stream直接架構串流還是架構登錄檔支援串流。 詳情請參見 更新串流 。

管理串流

去直播

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。

更新串流

創建之後可以用 update_stream 來更改串流。 通過 schema_config 以演化直接模式、 excluded_columns 改變刪除欄位,或兩者皆可。 不支援更新其他欄位。 改為建立新的串流。

更新串流會重新啟動其擷取管線,因此變更會生效。 通常幾分鐘內會恢復吞嚥。

演化出直接結構

對於使用直接結構的串流,將 a DirectSchemas 傳遞給 schema_config。 設定 payload_schema、 key_schema或兩者皆有。 你沒設定的那一面則保持不變。 結構登錄檔支援的串流會拒絕 schema_config 更新,必須透過登錄檔演化。

from databricks.feature_engineering.entities import DirectSchemas, SchemaConfig

stream = client.update_stream(
    name="my_catalog.my_schema.my_stream",
    schema_config=DirectSchemas(
        payload_schema=SchemaConfig(
            json_schema=(
                '{'
                '  "type": "object",'
                '  "properties": {'
                '    "user_id": {"type": "string"},'
                '    "amount": {"type": "number"},'
                '    "event_time": {"type": "string"},'
                '    "channel": {"type": "string"}'
                '  }'
                '}'
            )
        ),
    ),
)

結構描述更新必須具備回溯相容性,讓執行中的擷取管線能持續解碼現有記錄並寫入擷取資料表。 任何其他變更一律不予接受。

允許的節目內容取決於賽制:

  • JSON 與 Protobuf:新增可選欄位、移除欄位,並擴寬欄位型別(例如,將欄位型別從 int 擴寬為 bigint)。 Protobuf 也允許重新排序欄位。
  • Avro:只允許將intlong後方欄位擴寬並移除,該欄位的位元組無法讀取。 若要更自由地演化 Avro 架構,建議改用由 schema registry 支援的串流。

新增欄位會擴大擷取表的解碼 key 和 value 結構。 更新前寫入的資料列保持原始形狀,新增欄位則對於那些較早的資料列讀取為 NULL。 移除與變更類型僅對更新後擷取的紀錄生效。

變更排除欄位

將完整的一組新資料欄路徑傳遞至 excluded_columns,這會取代現有的集合。 傳送一個空清單([])以清除所有排除項目。 如需此行為的詳細資訊,請參閱 從串流中排除欄位。

stream = client.update_stream(
    name="my_catalog.my_schema.my_stream",
    excluded_columns=["value.user.email", "value.user.ssn"],
)

更改排除欄位只能往後進行。 新排除在外的欄位會停止寫入(顯示於 NULL),而新納入的欄位自此之後會開始填入資料,先前已寫入的資料列則維持不變。 為了防止新的欄位永遠不會被匯入:

  • Schema registry:先將欄位加入 excluded_columns,等待匯入管線重新啟動,然後在 registry 中註冊新的 schema 版本。
  • 直接結構:在同一呼叫中update_stream將欄位加入 和 schema_config 。excluded_columns

刪除串流

刪除串流同時也會刪除其擷取管線和擷取表。

Warning

任何參考已刪除串流的模型或特徵將無法存取底層串流資料。 如果你需要這些資料但不再需要串流,請先在刪除前建立一個擷取表的副本。

client.delete_stream(name="my_catalog.my_schema.my_stream")

範例筆記本

關於建立串流、定義串流功能並部署到服務端端的端對端範例,請參考以下筆記本:

串流功能檢視快速入門筆記本

拿筆記本