스트림 설정

중요합니다

이 기능은 공개 미리보기 단계에 있습니다. 작업 영역 관리자는 미리 보기 페이지에서 이 기능에 대한 액세스를 제어할 수 있습니다. Azure Databricks 미리 보기 관리를 참조하세요.

스트림은 Apache Kafka와 같은 외부 스트리밍 데이터 원본을 나타냅니다. 스트림은 연결 세부 정보, 인증, 스키마 및 수집 구성을 저장합니다. 스트림이 생성된 후에는 Feature View 정의에서 이를 참조하여 실시간 스트리밍 기능을 생성할 수 있습니다.

스트림에는 세 부분으로 구성된 이름(catalog.schema.stream_name)이 있습니다. Stream에 대한 액세스는 연결된 수집 테이블에 의해 제어됩니다. 자세한 내용은 수집 및 백필을 참조하세요.

요구 사항

  • Notebook 명령을 실행하는 경우: Databricks Runtime 17.0 ML 이상을 실행하는 서버리스 또는 클래식 컴퓨팅 클러스터입니다.
  • feature-engineering-client Python 패키지 버전 0.16.0 이상을 설치해야 합니다.

스트림 만들기

새 스트림을 만드는 데 사용합니다 create_stream() . 스트림에는 다음 네 가지 구성 요소가 필요합니다.

  • 원본 구성: 스트리밍 플랫폼(예: Kafka) 및 원본별 세부 정보(예: 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"
        ),
    ),
)

스트림 소스에 연결

스트리밍 기능을 정의하기 전에 Kafka broker에 대한 스트리밍 Lakeflow 파이프라인 연결을 연결하고 테스트합니다. 서버리스 컴퓨팅에서 스트리밍을 참조하고 Apache Kafka에 연결합니다.

AWS 관리 스트리밍(Amazon MSK)의 경우 Amazon MSK에 대한 서버리스 프라이빗 연결을 참조하세요. Kafka 인증 옵션에 대한 자세한 내용은 인증을 참조 하세요.

Authentication

Unity 카탈로그 연결을 사용하여 Kafka 클러스터에 인증합니다. 관리 인증에 권장되는 방법입니다. 연결을 만들려면 연결 만들기를 참조하세요.

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

직접 mTLS

직접 mTLS 인증을 위해 Unity Catalog 볼륨에 저장된 키스토어 및 트러스트스토어 파일을 제공하세요. 비밀번호는 Databricks 시크릿 스코프를 통해 참조되어야 합니다. 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"
        ),
    ),
)

SASL

SASL 인증(SASL/SCRAM 및 SASL/PLAIN 모두)은 미리 보기 중에 지원되지 않습니다.

구독 모드

구독 모드는 Stream에서 사용할 Kafka 토픽을 선택하는 방법을 지정합니다. 세 가지 모드가 지원됩니다.

모드 Description 예시
subscribe 토픽 이름의 쉼표로 구분된 목록 KafkaSubscriptionMode(subscribe="topic1,topic2")
subscribe_pattern Java 정규식 패턴과 일치하는 토픽 이름 KafkaSubscriptionMode(subscribe_pattern="events-.*")
assign 토픽 파티션 할당을 지정하는 JSON KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

스키마 구성

JSON 스키마 형식을 사용하여 메시지 키 및 값의 구조를 정의합니다. Kafka 원본의 payload_schema 경우 Kafka 메시지 값( value Kafka의 키-값 모델)에 해당하며 key_schema Kafka 메시지 키에 해당합니다. 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 수집 테이블에서 Stream에 대한 읽기 권한을 부여합니다.
  • 수집 테이블의 MANAGE는 삭제 권한을 부여합니다.

테이블 권한에 대한 자세한 내용은 테이블Unity 카탈로그 권한 참조를 참조하세요.

데이터 수집 파이프라인

스트림이 생성되면 Databricks는 Kafka 토픽에서 메시지를 지속적으로 읽어 Delta 테이블(즉, 수집 테이블)에 기록하는 관리형 수집 파이프라인을 시작합니다. 파이프라인은 최신 Kafka 오프셋에서 시작하여 지속적으로 실행되어 스트림을 만든 후에 도착하는 새 메시지만 캡처합니다. 이 수집 테이블은 스트리밍 기능을 사용하여 학습하는 데 사용됩니다. 스트림이 삭제되면 해당 수집 파이프라인 및 수집 테이블도 삭제됩니다.

수집 대상

스트림 ingestion_destination 데이터가 기록되는 세 부분으로 구성된 델타 테이블 이름을 지정합니다.

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

수집 테이블 스키마

수집 테이블에는 메타데이터 열과 함께 메시지 데이터가 포함됩니다.

Column Type Description
key 출발지(key_schema)에 따라 달라집니다 제공한 스키마에 따라 구조화된 Kafka 메시지 키입니다.
value 출발지(payload_schema)에 따라 달라집니다 제공한 스키마에 따라 구조화된 Kafka 메시지 값(페이로드)입니다.
stream_record_timestamp TIMESTAMP 기록의 타임스탬프입니다. forward-fill 데이터의 경우, 이는 Kafka 브로커 수집 타임스탬프입니다. 백필 데이터는 고객이 제공한 데이터입니다.
kafka_topic STRING 레코드가 소비된 Kafka 토픽입니다.
kafka_partition INT 레코드가 소비된 Kafka 파티션입니다.
kafka_offset LONG 해당 파티션 내 레코드의 카프카 오프셋입니다.
record_source STRING "stream"(라이브 Kafka 스트림에서 포워드 필) 또는 "backfill"(백필 소스에서) 둘 중 하나입니다.

백필 소스

정방향 채우기 파이프라인은 최신 Kafka 오프셋에서 시작되므로 스트림을 만들기 전에 존재했던 메시지를 캡처하지 않습니다. 학습용 과거 데이터 범위를 제공하려면 선택적 백필 소스를 구성하세요.

백필 소스가 구성되면 Databricks는 백필 행을 MERGE INTO와 함께 인제스트 테이블에 복사하는 일회성 record_source="backfill" 작업을 실행합니다. MERGE는 중첩 검사기가 백필 원본과 앞으로 채우기 스트림에 겹치는 타임스탬프가 있음을 확인한 후에만 실행됩니다( 백필과 라이브 스트림 데이터 간의 겹침 참조). 중첩 조건이 2일 이내에 충족되지 않더라도 무기한 차단을 방지하기 위해 MERGE는 계속 실행됩니다.

백필 테이블에는 UTC 시간대의 stream_record_timestamp 유형인 TIMESTAMP 열이 포함되어야 합니다. 다른 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")입니다.

MERGE는 백필의 최신 타임스탬프가 수집 테이블의 가장 빠른 타임스탬프를 1시간 이상 초과하면 진행됩니다. 이 겹침으로 인해 수집 테이블에 공백이 생기지 않습니다. 중첩 조건이 2일 이내에 충족되지 않더라도 무기한 차단을 방지하기 위해 MERGE는 계속 실행됩니다.

수집 파이프라인은 최신 Kafka 오프셋에서 시작되므로 스트림을 만든 후에 도착하는 메시지만 캡처합니다. 백필 원본에는 스트림 생성 시간뿐만 아니라 수집 시간 범위로 확장되는 데이터가 포함되어야 합니다.

예를 들어 오후 3시에 스트림을 만드는 경우 전달 채우기 파이프라인은 오후 3시 이후부터 메시지를 읽기 시작합니다. 겹침 검사를 충족하려면 백필 소스에 최소 오후 4:00까지의 타임스탬프가 포함된 데이터(포워드필 시작 시점보다 1시간 후)가 포함되어 있어야 합니다. 즉, 수집 테이블에 누락 구간이 없도록 오후 4시 이후에 백필 테이블을 업데이트해야 합니다.

Deduplication

deduplication_columns를 사용하여 백필과 포워드필 스트림 데이터 수집 중 중복 행을 식별할 열 경로를 지정합니다. 중첩 필드(예 "value.user_id": )에 점 표기법을 사용합니다.

데이터에 따라 중복 제거 열을 선택합니다.

  • 스트림의 각 레코드에 고유 식별자(예: value.transaction_id)가 포함된 경우 중복 제거에 해당 열을 사용합니다.
  • 백필 소스에 kafka_partitionkafka_offset 열이 포함되어 있으면 해당 열을 사용하여 각 레코드를 고유하게 식별합니다.
  • 중복 제거 열이 지정되지 않은 경우 기본 중복 제거 키는 , keyvalue.의 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")

예제 노트

Stream을 만들고, 스트리밍 기능을 정의하고, 서비스 엔드포인트에 배포하는 엔드 투 엔드 예제는 다음 Notebook을 참조하세요.

스트리밍 기능 보기 빠른 시작 노트북

노트북 받기