Настройка потока

Это важно

Эта функция доступна в общедоступной предварительной версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".

Stream представляет внешний источник данных потоковой передачи, например Apache Kafka. Потоки хранят сведения о подключении, данные аутентификации, схемы и конфигурацию загрузки данных. После создания потока вы можете ссылаться на него с помощью определений представления функций для создания функций потоковой передачи в режиме реального времени.

Потоки имеют три части имен (catalog.schema.stream_name). Доступ к потоку данных определяется связанной с ним таблицей приема данных. Подробнее см. в разделе Приём данных и обратное заполнение.

Требования

  • Для выполнения команд записной книжки: бессерверный или классический вычислительный кластер под управлением Databricks Runtime 17.0 ML или более поздней версии.
  • feature-engineering-client Необходимо установить пакет Python версии 0.16.0 или более поздней версии.

Создание потока

Используйте create_stream(), чтобы создать новый 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 и проверьте его подключение к конвейеру потоковой передачи Lakeflow. См. Потоковая передача в бессерверной вычислительной среде и Подключение к Apache Kafka.

Для управляемой потоковой передачи AWS (Amazon MSK) см. бессерверные частные подключения к Amazon MSK. Дополнительные сведения о параметрах проверки подлинности Kafka см. в разделе "Проверка подлинности".

Authentication

Используйте подключение Unity Catalog для аутентификации в кластере Kafka. Это рекомендуемый подход для управляемой проверки подлинности. Сведения о создании подключения см. в разделе "Создание подключения".

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

Direct mTLS

Для прямой проверки подлинности mTLS укажите хранилище ключей и файлы truststore, хранящиеся в томе каталога Unity, с паролями, на которые ссылаются области секретов Databricks. Дополнительные сведения о проверке подлинности SSL с помощью Kafka см. в статье "Использование 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 для использования. Поддерживаются три режима:

Режим Description Example
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 в модели key-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 Catalog.

Конвейер загрузки данных

При создании потока Databricks запускает управляемый конвейер приема данных, который непрерывно считывает сообщения из топика Kafka и записывает их в таблицу Delta (таблицу приема данных). Конвейер запускается с последнего смещения Kafka и работает непрерывно, захватывая только новые сообщения, поступающие после создания потока. Эта таблица приема используется для обучения с функциями потоковой передачи. При удалении потока также удаляются его конвейер приема и таблица приема.

Место назначения для приема данных

ingestion_destination указывает трёхсоставное имя таблицы Delta, в которую записываются потоковые данные.

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

Схема таблицы приема

Таблица приема данных содержит данные сообщения, а также столбцы метаданных:

Column Тип Description
key Зависит (от key_schema) Ключ сообщения Kafka, структурированный в соответствии с предоставленной схемой.
value Зависит (от payload_schema) Значение сообщения Kafka (полезные данные), структурированное в соответствии со схемой, которую вы предоставили.
stream_record_timestamp TIMESTAMP Временная метка записи. Для данных с дозаполнением это метка времени приема брокером Kafka. Для ретроспективных данных эти данные предоставляет клиент.
kafka_topic STRING Топик Kafka, из которого была получена запись.
kafka_partition INT Раздел Kafka, из которого была получена запись.
kafka_offset LONG Смещение записи Kafka в пределах её партиции.
record_source STRING Либо "stream" (дозаполнение из потока Kafka в реальном времени), либо "backfill" (из источника обратного заполнения).

Источник дозаполнения

Поскольку конвейер дозаполнения начинается с последнего смещения Kafka, он не захватывает сообщения, существовавшие до создания потока. Чтобы обеспечить охват исторических данных для обучения, настройте необязательный источник дозагрузки данных.

Если настроен источник для обратного заполнения, Databricks выполняет однократное задание MERGE INTO, которое копирует строки обратного заполнения в таблицу приема данных с помощью record_source="backfill". MERGE выполняется только после того, как средство проверки перекрытия подтвердит, что источник обратного заполнения и поток прямого заполнения имеют пересекающиеся временные метки (см. Перекрытие между данными обратного заполнения и данными потока в реальном времени). Если условие перекрытия не выполняется в течение 2 дней, операция MERGE всё равно выполняется, чтобы не допустить блокировки на неопределённый срок.

Таблица дозаполнения должна содержать столбец 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") в таблице загрузки.

Операция MERGE выполняется, когда самая поздняя метка времени дозагрузки исторических данных превышает самую раннюю метку времени таблицы загрузки как минимум на 1 час. Это перекрытие гарантирует отсутствие пробелов в таблице приема. Если условие перекрытия не выполняется в течение 2 дней, операция MERGE всё равно выполняется, чтобы не допустить блокировки на неопределённый срок.

Поскольку конвейер приёма данных начинается с последнего смещения Kafka, он захватывает только сообщения, поступающие после создания потока. Источник для дозаполнения должен содержать данные, охватывающие временной диапазон приёма данных, а не только период до времени создания потока данных.

Например, если вы создаете поток в 15:00, конвейер дозаполнения начинает считывать сообщения начиная с 15:00. Источник для обратного заполнения должен включать данные с метками времени как минимум до 16:00 (на 1 час позже начала прямого заполнения), чтобы пройти проверку перекрытия. Это означает, что вам следует обновить таблицу дозагрузки после 16:00, чтобы в таблице загрузки не было пропусков.

Дедупликация

Используйте deduplication_columns для указания путей к столбцам, чтобы определять дублирующиеся строки при приеме данных между данными backfill и потоком данных forward-fill. Используйте нотацию точек для вложенных полей (например, "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.

Удалите поток

При удалении потока также удаляется конвейер приема данных и таблица приема данных.

Предупреждение

Любые модели или функции, ссылающиеся на удаленный поток, больше не будут иметь доступа к базовым данным потока. Создайте копию таблицы приема данных перед удалением потока, если эти данные вам нужны, но сам поток больше не нужен.

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

Пример записной книжки

Полный пример создания Stream, определения функций потоковой передачи и развертывания в конечной точке обслуживания см. в следующей записной книжке:

Записная книжка быстрого запуска представлений функций потоковой передачи

Получите ноутбук