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

Это важно

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

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

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

Требования

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

Подключение к источникам потоков

Перед определением функций потоковой передачи подключитесь к брокеру Kafka и проверьте его подключение к конвейеру потоковой передачи Lakeflow. Feature Store опирается на бессерверную платформу SDP, а это означает, что вам потребуется механизм для подключения ваших классических вычислительных ресурсов (брокера или конечной точки) к бессерверным вычислительным ресурсам Databricks. Это достигается с помощью таких продуктов, как Private Link, или путём предоставления доступа к вашим классическим вычислительным ресурсам из общедоступного интернета.

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

Используйте create_stream(), чтобы создать новый Stream. Для потока требуется четыре компонента конфигурации:

  • Конфигурация исходного кода: указывает платформу для стриминга и специфичные для источника детали, такие как тематическая подписка для исходника Kafka.
  • Конфигурация подключения. Указывает, как подключиться и пройти проверку подлинности на платформе потоковой передачи, включая серверы начальной загрузки и учетные данные.
  • Конфигурация схемы: определяет структуру ключей и значений сообщений.
  • Конфигурация приема данных: указывает, откуда и каким образом поступают потоковые данные. См. «Прием данных и обратная загрузка» для получения подробной информации.

Сведения о source_config, специфичном для конкретного источника, и настройке подключения, а также полный пример create_stream() см. в Apache Kafka. Параметры схемы и приёма данных являются общими для всех источников.

Apache Kafka

Для трансляции с Apache Kafka используйте KafkaStreamConfig в качестве исходной конфигурации и подключение Unity Catalog для аутентификации. См. «Потоковая передача в бессерверной среде вычислений» и «Подключение к Apache 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

Режим подписки указывает, как 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]}')

Аутентификация Кафки

Используйте подключение Unity Catalog для аутентификации в кластере Kafka. Это рекомендуемый подход для управляемой проверки подлинности. Сведения о создании подключения см. в разделе "Создание подключения". Для создателя потока должен быть доступен USE CONNECTION в подключении. Любой пользователь, использующий материализованные функции, где Stream выступает в качестве источника, должен также иметь USE CONNECTION для этого подключения.

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

Соединение поддерживает как IAM (учетные данные сервиса), так и SASL-аутентификацию.

IAM (удостоверение службы)

Аутентифицироваться с помощью учетных данных Unity Catalog, например, чтобы подключиться к Amazon MSK с IAM. Сведения о создании учетных данных службы см. в разделе "Создание учетных данных службы". Задайте имя учетных данных сервиса с помощью опции credential:

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

Помимо USE CONNECTION подключения, на ACCESS нём должны быть идентификаторы, использующие учетные данные сервиса. Присвоите ACCESS указанный сервисный удостоверитель создателю Потока и любую личность, которая реализует функции с Stream. См . раздел "Предоставление разрешений на использование учетных данных службы для доступа к внешней облачной службе".

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 укажите хранилище ключей и файлы 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"
        ),
    ),
)

Конфигурация схемы

Определите структуру ключей и значений сообщений так, чтобы при загрузке данных и в определениях признаков можно было читать отдельные поля. Для источников Kafka payload_schema соответствует значению сообщения Kafka (value в модели key-value Kafka), а key_schema соответствует ключу сообщения Kafka. Необходимо указать по крайней мере одно из значений: payload_schema или key_schema.

Каждый SchemaConfig из них принимает один из трёх форматов, соответствующих тому, как исходный источник сериализирует свои сообщения: json_schema, avro_schema, или proto_schema. Если для ключа или полезных данных не указана схема, она рассматривается как простая строка.

Примеры кода в этом разделе используют схемы, объявленные непосредственно в коде с помощью DirectSchemas, при этом схема задаётся в виде строки. Для управления схемами с использованием внешнего реестра схем см. Реестр схем для подробностей.

Схема JSON

Укажите строку JSON Schema для 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

Укажите a ProtoSchemaSpecproto_schema с исходным текстом буферов.proto протокола и именем сообщения полезной нагрузки. Импорт 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-значение вместо того, чтобы привести к сбою потока.
  • Профсоюзы Avro. Объединение нескольких типов записей декодируется в структуру, содержащую по одному полю для каждого типа записи, при этом каждое поле получает имя соответствующей записи Avro.
  • Типы Протобуф. Беззнаковые целые числа декодируются в более широкий знаковый тип (например, uint32 в BIGINT и uint64 в DECIMAL(20,0)), поля enum-декодируются на имя строки, а типы скалярных оболочек (например StringValue , и Int32Value) декодируются в нулируемый столбец обёрнутого типа.

Реестр схем

Реестры схем хранят и используют схемы версий, которые используют производители и потребители потокового потока, обеспечивая соблюдение правил совместимости по мере развития этих схем. Когда настраивается внешний реестр схемы, Feature Store считывает схему из реестра и использует её для декодирования потокового сообщения. При использовании реестра схем не нужно объявлять схему непосредственно в Stream.

Поддержка реестра схем имеет следующие ограничения:

  • Поддерживается только для стримов Kafka.
  • Поддерживается только Confluent Schema Registry
  • Поддерживаются только форматы Avro и Protobuf . Чтобы читать JSON-сообщения, вместо этого объявьте схему в строке. См. схему JSON.
  • Каждый поток подключён ровно к одному субъекту Confluent для значения сообщения и к одному — для ключа сообщения (если он задан). Темы потока, содержащие несколько записей схемы, не являются поддерживаемой конфигурацией. Если ваш поток подключается к темам, содержащим несколько схем, записи, не соответствующие схеме указанного предмета, декодируются как null.

Подключитесь к реестру схемы

Укажите параметры подключения к реестру в качестве параметров подключения Kafka к Unity Catalog и храните секрет API реестра в Databricks secret scope. Идентичность потока 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>')
)

Установите для параметра schema_registry_api_secret в подключении Kafka и для ссылки на область секретов в потоке один и тот же секрет.

Создайте поток, использующий реестр схем

Передайте SchemaRegistryConfig в качестве schema_config. Укажите секрет API реестра с помощью api_secret_ref, а также определите субъект и формат с помощью 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 в соответствии с тем, как сериализуется раздел.

Развитие схемы

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

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

О том, как конвейер обрабатывает записи, не соответствующие используемой схеме, см. раздел «Декодирование данных с использованием схем».

Поглощение и заполнение

Параметр ingestion_config настраивает запись и хранение потоковых данных для обучения и обслуживания.

Доступ к потоку регулируется таблицей загрузки данных:

  • SELECT В таблице приема предоставляется доступ на чтение к Stream.
  • MANAGE предоставляет права на удаление в таблице загрузки.

Дополнительные сведения о правах доступа к таблицам см. в разделах Таблица и Справочник по привилегиям Unity Catalog.

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

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

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

ingestion_destination указывает трёхсоставное имя таблицы 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, которое копирует строки обратного заполнения в таблицу приема данных с помощью record_source="backfill". MERGE выполняется только после того, как средство проверки перекрытия подтвердит, что источник обратного заполнения и поток прямого заполнения имеют пересекающиеся временные метки (см. Перекрытие между данными обратного заполнения и данными потока в реальном времени). Если условие перекрытия не выполняется в течение 2 дней, операция MERGE всё равно выполняется, чтобы не допустить блокировки на неопределённый срок.

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

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

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

Например, если вы создаете поток в 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, определения функций потоковой передачи и развертывания в конечной точке обслуживания см. в следующей записной книжке:

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

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