Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Это важно
Эта функция доступна в общедоступной предварительной версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями 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 (рекомендуется)
Используйте подключение 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, определения функций потоковой передачи и развертывания в конечной точке обслуживания см. в следующей записной книжке: