ストリームを設定する

Important

この機能は パブリック プレビュー段階です。 ワークスペース管理者は、[ プレビュー] ページからこの機能へのアクセスを制御できます。 Manage Azure Databricks プレビューを参照してください。

Stream は、Apache Kafka などの外部ストリーミング データ ソースを表します。 ストリームには、接続の詳細、認証、スキーマ、およびインジェスト構成が格納されます。 ストリームが作成されたら、 フィーチャ ビュー 定義を使用してストリームを参照し、リアルタイム ストリーミング 機能を作成できます。

ストリームには 3 部構成の名前 (catalog.schema.stream_name) があります。 ストリームへのアクセスは、関連するインジェスト テーブルによって管理されます。 詳細は 「インジェストとバックフィル」 を参照してください。

必要条件

  • ノートブック コマンドを実行する場合: サーバーレスまたは Databricks Runtime 17.0 ML 以上を実行しているクラシック コンピューティング クラスター。
  • feature-engineering-client Python パッケージのバージョン 0.18.0 以降がインストールされている必要があります。

ストリーム ソースへの接続

ストリーミング機能を定義する前に、Kafka ブローカーにストリーミング Lakeflow パイプライン接続を接続してテストします。 Feature StoreはサーバーレスSDPに依存しているため、従来のコンピュート(ブローカーやエンドポイント)をDatabricksのサーバーレスコンピュートに接続する仕組みが必要です。 これはprivatelinkのような製品や、従来のコンピュートをパブリックインターネットからアクセスできるようにすることで実現されます。

ストリームの作成

create_stream()を使用して新しいストリームを作成します。 Stream には、次の 4 つの構成コンポーネントが必要です。

  • ソース設定:ストリーミングプラットフォームやソース固有の詳細(例えばKafkaソースのトピック購読など)を指定します。
  • 接続構成: ブートストラップ サーバーと資格情報など、ストリーミング プラットフォームに接続して認証する方法を指定します。
  • スキーマ構成: メッセージ キーと値の構造を定義します。
  • インジェスト構成: ストリーム データを取り込む場所と方法を指定します。 詳細は 「インジェストとバックフィル」 を参照してください。

ソース固有の source_config と接続設定、そして完全な create_stream() 例については、 Apache Kafkaを参照してください。 スキーマとインジェス(取り込み)オプションはソース間で共有されます。

アパッチ・カフカ

Apache Kafkaからストリーミングするには、 KafkaStreamConfig をソース設定に使い、認証にはUnity Catalog接続を使います。 サーバーレス コンピューティングでのストリーミングおよび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 トピックを選択する方法を指定します。 次の 3 つのモードがサポートされています。

Mode Description 例
subscribe トピック名のコンマ区切りリスト KafkaSubscriptionMode(subscribe="topic1,topic2")
subscribe_pattern Javaの正規表現でトピック名を照合する KafkaSubscriptionMode(subscribe_pattern="events-.*")
assign トピック パーティションの割り当てを指定する JSON KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

カフカ認証

Unity カタログ接続を使用して、Kafka クラスターに対する認証を行います。 これは、マネージド認証に推奨される方法です。 接続を作成するには、「接続の 作成」を参照してください。 ストリームの作成者は、その接続で USE CONNECTION を持っている必要があります。 Stream をソースとして特徴量をマテリアライズするユーザーには、その接続に対する USE CONNECTION も必要です。

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

接続はIAM(サービス認証)とSASL認証の両方をサポートしています。

IAM(サービス認証情報)

例えば、IAM を使ってAmazon MSKに接続するためにUnity Catalogのサービス認証情報で認証してください。 サービス資格情報を作成するには、「 サービス資格情報の作成」を参照してください。 サービス認証名を credential オプションで設定します:

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

接続には USE CONNECTION が必要であり、さらに、サービス認証情報を使用する ID には 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 認証の場合は、Databricks シークレット スコープを介して参照されるパスワードを使用して、Unity カタログ ボリュームに格納されているキーストア ファイルとトラストストア ファイルを提供します。 Kafka での SSL 認証の詳細については、「SSL を使用して kafka にAzure Databricks接続する」を参照してください。

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

スキーマ構成

インジェスや機能定義が個々のフィールドを読み取れるように、メッセージキーや値の構造を定義します。 Kafka ソースの場合、 payload_schema は Kafka メッセージ値 (Kafka のキー値モデルの value ) に対応し、 key_schema は Kafka メッセージ キーに対応します。 少なくとも 1 つの payload_schema または key_schema を指定する必要があります。

各 SchemaConfig は、送信元がメッセージをシリアライズする方法に合わせて、 json_schema、 avro_schema、または proto_schemaの3つの形式のいずれかを受け入れます。 キーまたはペイロードにスキーマが指定されていない場合は、単純な文字列として扱われます。

このセクションのコード例は、 DirectSchemasとインラインで宣言されたスキーマを使用し、スキーマは文字列として提供されています。 外部スキーマレジストリを使ってスキーマを管理するには、スキー マレジストリ の詳細をご覧ください。

JSON スキーマ

json_schemaにJSONスキーマ文字列を提供します。

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 スキーマ

Protocol Buffers.proto のソース テキストとペイロードのメッセージ名を指定した ProtoSchemaSpec を proto_schema に渡します。 ProtoSchemaSpec から databricks.feature_engineering.entities をインポートします。

message_name は、package テキストで宣言された .proto を含む完全修飾メッセージ名である必要があります(たとえば、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 モードを用い、スキーマに合わないレコードはストリームに失敗するのではなく、null値にデコードされます。
  • アブロ・ユニオン。 複数のレコード型のユニオンは、各レコード型につき1フィールドを持つ構造体にデコードされ、各フィールドには対応するAvroレコード名が付けられます。
  • プロトブフタイプ。 符号なし整数はより広い符号付き型にデコードされます(例: uint32 が BIGINT 、 uint64 が DECIMAL(20,0)に)、enumフィールドは文字列名に復号し、スカラーラッパー型(例: StringValue や Int32Value)はラップ型のnullable列に復号します。

スキーマ レジストリ

スキーマレジストリーはストリーミング制作者や消費者が使用するバージョンスキーマを保存し、それらのスキーマが進化するにつれて互換性ルールを強制します。 外部スキーマレジストリが設定されると、Feature Storeはレジストリからスキーマを読み取り、それを使ってストリーミングメッセージをデコードします。 スキーマレジストリを使う際にストリーム上でスキーマをインラインで宣言する必要はありません。

スキーマレジストリのサポートには以下の制限があります:

スキーマレジストリに接続する

Kafka Unity Catalog接続のオプションとしてレジストリ接続の詳細を提供し、レジストリAPIの秘密をDatabricks の秘密スコープに保存します。 ストリームの実行時アイデンティティは秘密スコープに対して READ 許可を持たなければなりません。なぜなら、インジェスションパイプラインは実行時に秘密を読み取るからです。 接続の作成および設定方法については、「 接続を作成」をご覧ください。

認証に使われる接続にschema_registry_url、schema_registry_api_key、schema_registry_api_secretオプションを追加してください。 以下の例は、Unity Catalogのサービス認証情報でブローカーに認証し、APIキーでレジストリに認証するKafka接続を作成します。

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接続の schema_registry_api_secret オプションとStreamの シークレットスコー プ参照の両方を同じシークレットに設定してください。

スキーマレジストリを使用するストリームを作成します

SchemaRegistryConfigとしてschema_configを渡します。 レジストリAPIシークレットを api_secret_refで参照し、件名とフォーマットはメッセージ値として payload_schema_locator 、メッセージキーには key_schema_locator で識別します。 少なくとも1つのロケーターが用意されなければなりません。

ここで、スキーマ構成 セクションにある直接的なスキーマの例との違いに注意してください。 スキーマレジストリを使用する場合、ストリームから 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 に設定してください。

スキーマの進化

インジェストパイプラインは、起動時にサブジェクトの現在のスキーマを特定します。 スキーマレジストリで対象の新しい後方互換スキーマバージョンを登録すると、実行中のパイプラインは開始時のバージョンを引き続き使用します。

スキーマ レジストリを利用したストリームの場合、取り込みパイプラインは数時間ごとに自動的に再起動します。 再起動のたびに、対象の最新のスキーマ バージョンが読み込まれ、新規または変更されたフィールドが取り込みテーブルに表示されます。

スキーマ レジストリの代わりにダイレクト スキーマを使用するストリームは、update_stream を使用してスキーマを進化させます。 詳細は、 ストリームの更新を参照してください。

パイプラインが現在使用しているスキーマと一致しないレコードをどのように扱うかについては、「 スキーマを使ったデータのデコード」を参照してください。

レコードをタイプでフィルターする

ストリームは、キーと値のスキーマ (指定されている場合) を 1 つ使用して各レコードをデコードします。これは、それらを直接指定する場合でも、スキーマ レジストリを使用する場合でも同様です。 トピックには複数の種類のレコードが含まれる場合があり、ストリームは複数のトピックを購読できるため、record_type_filter を使用して、そのトピック内のどのレコードをこのストリームに含めるかを選択してください。

例えば、value.event_type = 'transaction' などのドット表記でデコードされたフィールドを参照するSQL式を指定します。 フィルタに一致しないレコードは無視されます。 これらは取り込みテーブルには書き込まれず、マテリアライゼーションでも使用されません。 他のレコード タイプ用のストリームを作成するには、record_type_filter を変更して別のストリームを作成します。

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 インジェスト テーブルで Stream への読み取りアクセスを許可します。
  • MANAGE インジェスト テーブルで削除アクセスを許可します。

テーブル権限の詳細については、「 Table および Unity Catalog 権限リファレンス」を参照してください。

インジェスト パイプライン

ストリームが作成されると、Databricksは管理型インジェスションパイプラインを開始し、ソースストリームからメッセージを継続的に読み込み、デルタテーブル(インジェスションテーブル)に書き込みます。 パイプラインはソースの最新の位置から始まり、ストリーム作成後に到着した新しいメッセージのみを継続的にキャプチャします。 このインジェスト テーブルは、 ストリーミング機能を使用したトレーニングに使用されます。 ストリームが削除されると、そのインジェスト パイプラインとインジェスト テーブルも削除されます。

取り込み先

ingestion_destinationは、ストリーム データが書き込まれる 3 部構成の Delta テーブル名を指定します。

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

インジェスト テーブル スキーマ

インジェスションテーブルにはメッセージデータとメタデータの列が含まれています。 共通の列はすべての源に存在します。 kafka_* 列はカフカ川にのみ存在します。

Column タイプ Source Description
key ( key_schemaによって異なります) Common メッセージキーは、あなたが提供したスキーマに基づいて構成されています。
value ( payload_schemaによって異なります) Common メッセージ値(ペイロード)は、あなたが提供したスキーマに基づいて構造化されています。
stream_record_timestamp TIMESTAMP Common レコードのタイムスタンプ。 フォワードフィルデータの場合、これはソースインジェストタイムスタンプです。 バックフィル データの場合、これは顧客が提供します。
record_source STRING Common "stream"(ライブストリームからのフォワードフィル)か"backfill"(バックフィルソースからの)のいずれかです。
kafka_topic STRING Kafka レコードが使用された Kafka トピック。
kafka_partition INT Kafka レコードが使用された Kafka パーティション。
kafka_offset LONG Kafka そのパーティション内における、そのレコードの Kafka オフセット。

補完の供給元

フォワードフィルパイプラインは送信元の最新の位置から始まるため、ストリーム作成前に存在していたメッセージはキャプチャしません。 トレーニング用の履歴データを確保するために、オプションのバックフィル ソースを設定します。

バックフィル ソースを構成すると、Databricks は、MERGE INTOを使用してインジェスト テーブルにバックフィル行をコピーする 1 回限りのrecord_source="backfill" ジョブを実行します。 MERGE は、バックフィル ソースと前方フィル ストリームのタイムスタンプが重複していることをオーバーラップ チェッカーが確認した後にのみ実行されます (「 バックフィルデータとライブ ストリーム データの重複」を参照)。 重複条件が 2 日以内に満たされない場合は、MERGE が実行され、ブロックが無期限に回避されます。

バックフィル テーブルには、UTC タイムゾーンでstream_record_timestamp型のTIMESTAMP列を含める必要があります。 他のメタデータ列は、バックフィルソースに存在する場合はそのまま引き継がれ、存在しない場合は 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 を実行する前に、重複チェックによって 2 つのテーブルのタイムスタンプが比較されます。

  • バックフィルの最大値: バックフィル ソース内の最大 stream_record_timestamp 。
  • 取り込みの最小値: インジェスト テーブル内の行数stream_record_timestamp(record_source="stream")の最小値。

MERGE は、バックフィルの最新のタイムスタンプがインジェスト テーブルの最も古いタイムスタンプを少なくとも 1 時間超えると続行されます。 この重複により、インジェスト テーブルにギャップがなくなります。 重複条件が 2 日以内に満たされない場合は、MERGE が実行され、ブロックが無期限に回避されます。

インジェスションパイプラインはソースの最新の位置から始まるため、ストリーム作成後に到着したメッセージのみをキャプチャします。 バックフィル ソースには、ストリームの作成時間だけでなく、インジェスト時間範囲まで拡張されるデータが含まれている必要があります。

たとえば、午後 3 時にストリームを作成した場合、前方入力パイプラインは午後 3 時以降にメッセージの読み取りを開始します。 バックフィル ソースには、重複チェックを満たすために、少なくとも午後 4:00 (前方入力開始から 1 時間) までのタイムスタンプを含むデータを含める必要があります。 つまり、インジェスト テーブルにギャップがないように、午後 4 時以降にバックフィル テーブルを更新する必要があります。

重複除去

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

コスト配分

ストリームのマネージド取り込みコストを割り当てるために、budget_policy_id に IngestionConfig と tags を設定します。 Azure Databricks は、Stream の作成時に、それらを取り込み用の Lakeflow パイプラインおよびそのフォワードフィル ジョブとバックフィル ジョブに適用します。

例えば、タグの制限や属性付いた支出のクエリ方法については、「 タグとサーバーレス使用ポリシーによる属性コスト」を参照してください。

ストリームから列を除外する

excluded_columns を使うと、取り込みたくない特定の列を Stream から除外できます。 除外された列は、取り込みテーブルには書き込まれず、フィーチャとして参照したり、トレーニングに使用したりすることはできません。

各列は、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 が失敗します。 スキーマ レジストリを使用する場合、その列がまだ存在しない段階で除外することができます。 除外された列は、重複排除列にもできません。これは、重複排除列は重複する行を特定するために必要とされるためです。 除外された列を参照するフィーチャー (エンティティ、時系列、入力など) は、作成に失敗します。

Stream の除外列は、作成後に update_stream を使用して、直接スキーマの Stream とスキーマレジストリに裏付けられた 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 を指定して削除される列を変更したり、あるいはその両方を行ったりします。 他のフィールドの更新はサポートされていません。 その代わりに、新しいストリームを作成してください。

ストリームを更新すると、その変更が反映されるよう、取り込みパイプラインが再起動されます。 取り込みは通常、数分以内に再開します。

ダイレクトスキーマを進化させる

ダイレクト スキーマを使用するストリームには、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: int から long への幅の拡張のみを許可し、後続のフィールドで読み取られないバイトを持つ末尾のフィールドを削除します。 Avro スキーマをより柔軟に発展させるには、代わりにスキーマ レジストリを基盤とするストリームを使用してください。

フィールドを追加すると、取り込みテーブルのデコード済み 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 に表示されます)、新たに追加された列には今後データが書き込まれていきますが、以前に書き込まれた行はそのまま残されます。 新しい列が取り込まれるのを完全に防ぐ方法:

  • スキーマ レジストリ: まず excluded_columns に列を追加し、取り込みパイプラインの再起動を待ってから、新しいスキーマ バージョンをレジストリに登録します。«»
  • ダイレクト スキーマ: 同じexcluded_columns呼び出しでupdate_streamとschema_configに列を追加します。

ストリームを削除する

ストリームを削除すると、そのインジェスト パイプラインとインジェスト テーブルも削除されます。

Warning

削除されたストリームを参照するモデルまたは機能は、基になるストリーム データにアクセスできなくなります。 このデータが必要であってもストリームが不要になった場合は、削除前にインジェスト テーブルのコピーを作成します。

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

ノートブックの例

Stream を作成し、ストリーミング機能を定義し、サービス エンドポイントにデプロイするエンド ツー エンドの例については、次のノートブックを参照してください。

ストリーミング機能ビューのクイック スタート ノートブック

ノートブックを入手