Pular para o conteúdo principal

Configurar uma transmissão

info

Visualização

Esse recurso está em Prévia Pública. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Pré-visualizações . Consulte Gerenciar prévias do Databricks.

Uma Transmissão representa uma fonte de dados de transmissão externa, como Apache Kafka. As transmissões armazenam detalhes de conexão, autenticação, esquemas e configuração de ingestão. Depois que uma transmissão é criada, ela pode ser referenciada usando definições de Feature View para criar recursos de transmissão em tempo real.

As transmissões têm nomes de três partes (catalog.schema.stream_name). O acesso a uma transmissão é regido pela tabela de ingestão associada. Consulte Ingestão e preenchimento retroativo para obter detalhes.

Requisitos​

  • Para executar comandos do Notebook: serverless ou um cluster de compute clássico executando Databricks Runtime 17.0 ML ou acima.
  • O pacote feature-engineering-client Python versão 0.18.0 ou acima deve estar instalado.

Conectando-se a fontes de transmissão​

Antes de definir os recursos de transmissão, conecte e teste uma conexão de LakeFlow Pipelines de transmissão com seu broker do Kafka. A Feature Store depende do SDP serverless, o que significa que você precisará de um mecanismo para conectar seu compute clássico (broker ou endpoint) ao compute serverless do Databricks. Isso é feito por meio de produtos como o privatelink ou permitindo que seu compute clássico seja acessível a partir da internet pública.

Criar uma Transmissão​

Use create_stream() para criar uma nova transmissão. Uma transmissão requer quatro componentes de configuração:

  • Source config : Especifica a plataforma de transmissão e detalhes específicos da origem, como a inscrição em um tópico para uma origem Kafka.
  • **Configuração da conexão**: Especifica como conectar e autenticar-se na plataforma de transmissão, incluindo servidores de bootstrap e credenciais.
  • Configuração de esquema: Define a estrutura das key e dos valores da mensagem.
  • Ingestion config : Especifica onde e como os dados de transmissão são ingeridos. Consulte Ingestion and backfill para obter detalhes.

For the source-specific source_config and connection setup, along with a complete create_stream() example, see Apache Kafka. The schema and ingestion options are shared across sources.

Apache Kafka​

Para transmissão do Apache Kafka, use KafkaStreamConfig como a configuração de origem e uma conexão do Unity Catalog para autenticação. Consulte Transmissão em compute serverless e Conectar ao Apache Kafka para obter conectividade com o Kafka.

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

Modos de inscrição do Kafka​

O modo de inscrição especifica como a Transmissão seleciona tópicos do Kafka para consumir. Três modos são suportados:

Mode

Descrição

Exemplo

subscribe

Lista de nomes de tópicos separados por vírgulas.

KafkaSubscriptionMode(subscribe="topic1,topic2")

subscribe_pattern

Padrão Java regex correspondente aos nomes de tópico

KafkaSubscriptionMode(subscribe_pattern="events-.*")

assign

JSON especificando as atribuições de tópicos-partições

KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Mode

Descrição

Exemplo

subscribe

Lista de nomes de tópicos separados por vírgulas.

KafkaSubscriptionMode(subscribe="topic1,topic2")

subscribe_pattern

Padrão Java regex correspondente aos nomes de tópico

KafkaSubscriptionMode(subscribe_pattern="events-.*")

assign

JSON especificando as atribuições de tópicos-partições

KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Autenticação do Kafka​

Conexão do Unity Catalog (recomendada)​

Use uma conexão do Unity Catalog para autenticar no seu cluster Kafka. Esta é a abordagem recomendada para autenticação gerenciada. Para criar uma conexão, consulte Criar uma conexão. O criador da transmissão deve ter USE CONNECTION na conexão. Qualquer usuário que materialize recursos com a transmissão como fonte também deve ter USE CONNECTION na conexão.

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

A conexão oferece suporte à autenticação IAM (credencial de serviço) e SASL.

IAM (credencial de serviço)​

Autentique com uma credencial de serviço do Unity Catalog, por exemplo, para se conectar ao Amazon MSK com IAM. Para criar uma credencial de serviço, consulte Criar credenciais de serviço. Defina o nome da credencial de serviço com a opção credential:

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

Além de USE CONNECTION na conexão, as identidades que usam a credencial de serviço precisam de ACCESS nela. Conceda ACCESS na credencial de serviço referenciada ao criador da transmissão e a qualquer identidade que materialize recursos com a transmissão. Consulte Conceder permissões para usar uma credencial de serviço para acessar um serviço de cloud externo.

SASL​

A autenticação SASL usa um nome de usuário e senha. Defina sasl_mechanism como um dos seguintes:

  • PLAIN
  • SCRAM-SHA-256
  • SCRAM-SHA-512

Forneça as credenciais com as opções user e password. A conexão armazena essas credenciais de forma segura.

O exemplo a seguir usa SASL/SCRAM. Para SASL/PLAIN, defina sasl_mechanism como PLAIN.

SQL
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 Direto​

Para autenticação mTLS direta, forneça arquivos de keystore e truststore armazenados em um volume do Unity Catalog, com senhas referenciadas por meio dos Secret Scopes do Databricks. Para obter mais informação sobre a autenticação SSL com Kafka, consulte Usar SSL para conectar o Databricks ao Kafka.

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

Configuração do esquema​

Defina a estrutura das keys e valores das mensagens para que as definições de ingestão e de recurso possam ler campos individuais. Para fontes Kafka, payload_schema corresponde ao valor da mensagem Kafka (o value no modelo key-value do Kafka) e key_schema corresponde à key da mensagem Kafka. Pelo menos um de payload_schema ou key_schema deve ser fornecido.

Cada SchemaConfig aceita um dos três formatos, correspondendo à forma como a origem serializa suas mensagens: json_schema, avro_schema ou proto_schema. Se nenhum esquema for fornecido para uma key ou payload, ele será tratado como uma string simples.

Os exemplos de código nesta seção usam esquemas declarados em linha com DirectSchemas, onde o esquema é fornecido como uma string. Para gerenciar esquemas usando um registro de esquema externo, consulte Registro de esquema para obter detalhes.

Esquema JSON​

Forneça uma string de Esquema JSON para json_schema.

Python
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"}'
),
)

Esquema Avro​

Forneça uma string de esquema Avro para avro_schema. Tipos lógicos Avro são suportados, incluindo timestamp-millis, date e decimal.

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

Esquema Protobuf​

Forneça um ProtoSchemaSpec para proto_schema com o texto de origem .proto de Protocol Buffers e o nome da mensagem de payload. Importar ProtoSchemaSpec de databricks.feature_engineering.entities.

message_name deve ser o nome da mensagem totalmente qualificado, incluindo o package declarado no texto .proto (por exemplo, com.example.Event, não Event). Ambas as sintaxes proto2 e proto3 são suportadas.

google.protobuf.Timestamp e os tipos de wrapper escalar (StringValue, Int32Value, etc.) são compatíveis, e suas importações são resolvidas automaticamente. Outros tipos conhecidos, como Duration, Struct e Any, são rejeitados; em vez disso, codifique esses valores como um escalar ou mensagem compatível. Os tipos escalares fixed32 e fixed64 e map com key que não são strings também não são compatíveis.

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

Decodificação de dados usando esquemas​

O Databricks decodifica cada mensagem com as funções from_json, from_avro e from_protobuf do Spark. Os seguintes comportamentos se aplicam independentemente de você declarar o esquema em linha ou resolvê-lo a partir de um registro de esquema:

  • Registros malformados. A decodificação usa o modo PERMISSIVE, portanto, um registro que não corresponde ao seu esquema é decodificado como um valor nulo em vez de falhar na transmissão.
  • Uniões Avro. Uma união de vários tipos de registro é decodificada em uma struct com um campo por tipo de registro, cada um nomeado de acordo com seu registro Avro.
  • Tipos Protobuf. Inteiros sem sinal são decodificados para um tipo com sinal mais amplo (por exemplo, uint32 para BIGINT e uint64 para DECIMAL(20,0)), campos enum são decodificados para seu nome de string e tipos de wrapper escalar (por exemplo, StringValue e Int32Value) são decodificados para uma coluna anulável do tipo encapsulado.

Registro de esquema​

Os registros de esquema armazenam e versionam esquemas que produtores e consumidores de transmissão usam, aplicando regras de compatibilidade à medida que esses esquemas evoluem. Quando um registro de esquema externo é configurado, a Feature Store lê o esquema do registro e o utiliza para decodificar a mensagem de transmissão. O senhor não declara o esquema em linha na transmissão ao usar um registro de esquema.

O suporte ao registro de esquema tem as seguintes limitações:

  • Compatível apenas com transmissões Kafka.
  • Somente o Confluent Schema Registry é suportado
  • Somente os formatos Avro e Protobuf são suportados. Para read.json messages, declare o esquema inline. Consulte esquema JSON.
  • Cada transmissão está conectada a exatamente um assunto Confluent para o valor da mensagem, e um para a key da mensagem (se fornecida). Tópicos de transmissão que contêm registros de vários esquemas não são uma configuração compatível. Se sua transmissão se conectar a tópicos que contêm vários esquemas, os registros que não corresponderem ao esquema do assunto especificado serão decodificados como nulos.

Conectar a um registro de esquema​

Forneça os detalhes da conexão do registro como opções na conexão do Kafka Unity Catalog e armazene o segredo da API do registro em um Secret Scope do Databricks. A identidade de execução (run-as) da transmissão deve ter a permissão READ no Secret Scope, pois o pipeline de ingestão lê o segredo em Runtime. Para saber como criar e configurar uma conexão, consulte Criar uma conexão.

Adicione as opções schema_registry_url, schema_registry_api_key e schema_registry_api_secret à conexão usada para autenticação. O exemplo a seguir cria uma conexão Kafka que autentica no broker com uma credencial de serviço do Unity Catalog e no registro com uma key de API:

SQL
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>')
)

Defina a opção schema_registry_api_secret na conexão Kafka e a referência ao Secret Scope na transmissão para o mesmo segredo.

Criar uma transmissão que usa um registro de esquema​

Passe um SchemaRegistryConfig como o schema_config. Referencie o segredo da API do registro com api_secret_ref e identifique o assunto e o formato com payload_schema_locator para o valor da mensagem, ou key_schema_locator para a key da mensagem. Pelo menos um localizador deve ser fornecido.

Observe as diferenças aqui em comparação com os exemplos de esquema direto na seção Configuração do esquema. Ao usar um registro de esquema, você não fornece o esquema em linha na transmissão para schema_config. Em vez disso, você especifica um SchemaRegistryConfig que identifica o esquema no registro.

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

Um assunto (subject) Confluent é o escopo nomeado sob o qual o histórico de versão de um esquema é registrado e a compatibilidade é imposta. Defina subject como o nome do escopo relevante, que é comumente determinado a partir da estratégia de nome de assunto:

  • TopicNameStrategy (default, deriva o assunto do nome do tópico): <topic>-value para o valor e <topic>-key para a chave. Por exemplo, o esquema de valor para o tópico transactions usa o assunto transactions-value.
  • RecordNameStrategy (deriva o assunto do nome do registro do esquema, independentemente do tópico): o nome do registro totalmente qualificado, como com.example.Payment. Este é o namespace e o nome do registro para Avro, ou o pacote e o nome da mensagem para Protobuf.
  • TopicRecordNameStrategy (combina os nomes de tópico e de registro): <topic>-<fully-qualified-record-name>, como transactions-com.example.Payment.

format é obrigatório. Defina-o como SchemaLocatorFormat.FORMAT_AVRO ou SchemaLocatorFormat.FORMAT_PROTOBUF para corresponder à forma como o tópico é serializado.

evolução do esquema​

O pipeline de ingestão resolve o esquema atual do assunto quando ele começa. Quando você registra uma nova versão de esquema compatível com versões anteriores no assunto no registro de esquema, o pipeline em execução continua a usar a versão com a qual começou.

Para transmissões apoiadas por um registro de esquema, o pipeline de ingestão é reiniciado automaticamente a cada duas horas, aproximadamente. A cada reinicialização, ele obtém a versão mais recente do esquema do assunto, e os campos novos ou alterados aparecem na tabela de ingestão.

Transmissões que usam esquemas diretos em vez de um registro de esquema evoluem seu esquema com update_stream. Consulte Atualizar uma transmissão.

Para saber como o pipeline lida com registros que não correspondem ao esquema que ele está usando atualmente, consulte Decodificando dados usando esquemas.

Filtrar registros por tipo​

Uma transmissão decodifica cada registro usando um único esquema de key e valor (se fornecido), independentemente de você especificá-los diretamente ou usar um registro de esquema. Como um tópico pode transportar mais de um tipo de registro e as transmissões podem assinar vários tópicos, use record_type_filter para selecionar quais registros do tópico pertencem a esta transmissão.

Forneça uma expressão SQL que faça referência a campos decodificados com notação de ponto, por exemplo, value.event_type = 'transaction'. Os registros que não correspondem ao filtro são ignorados. Eles não são gravados na tabela de ingestão e não são usados na materialização. Para criar uma transmissão para outros tipos de registro, crie uma transmissão separada com um record_type_filter diferente.

Python
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
record_type_filter="value.event_type = 'transaction'",
)

Even without record_type_filter, decoding never fails the transmissão. A record that doesn't match the configured schema is decoded permissively. The records decode in one of the following ways:

  • Em uma linha com valores de NULL para os campos que o esquema espera, mas que o registro omite (JSON, Avro e Protobuf).
  • Em uma linha contendo valores que pertencem a um tipo de registro diferente (somente Avro e Protobuf).

Para identificar se uma linha na tabela de ingestão pertence ao tipo de registro esperado, use uma das seguintes verificações:

  • Verifique se um campo é igual a um valor esperado; por exemplo, value.event_type = 'transaction' (preferencial para Avro e Protobuf).
  • Verifique se um campo é não-NULL, por exemplo, value.activity_id IS NOT NULL.

O uso de record_type_filter com transmissões separados é recomendado quando os esquemas diferem substancialmente entre os tipos de registro no tópico ou quando você deseja governar o acesso a cada tipo de registro de forma independente. Para manter os custos, o Databricks recomenda que você mantenha um pequeno número de transmissões, já que cada transmissão tem um pipeline de ingestão e uma tabela de ingestão separados. Cada transmissão também usa um compute separado no momento da materialização. Você pode usar filtros específicos de recursos para a materialização.

record_type_filter difere de filter_condition de um recurso. record_type_filter é definido na transmissão e controla quais registros são ingeridos e disponibilizados para todos os recursos que usam a transmissão como fonte, enquanto filter_condition é definido em um recurso individual e filtra as linhas antes da agregação. Consulte Condições de filtro em fontes de transmissão para obter mais detalhes sobre filter_condition.

Ingestão e preenchimento retroativo​

O parâmetro ingestion_config configura como os dados de transmissão são capturados e armazenados para treinamento e disponibilização.

O acesso a uma transmissão é governado pela tabela de ingestão:

  • SELECT na tabela de ingestão concede acesso de leitura à Transmissão.
  • MANAGE na tabela de ingestão concede acesso de exclusão.

Para mais informações sobre privilégios de tabela, consulte Tabela e referência de privilégios do Unity Catalog.

Pipeline de ingestão​

Quando uma transmissão é criada, o Databricks inicia um pipeline de ingestão gerenciado que lê continuamente as mensagens da transmissão de origem e as grava em uma tabela Delta (a tabela de ingestão). O pipeline começa a partir da posição mais recente na origem e é executado continuamente, capturando apenas novas mensagens que chegam após a criação da transmissão. Esta tabela de ingestão é usada para treinamento com recursos de transmissão. Quando uma transmissão é excluída, seu pipeline de ingestão e sua tabela de ingestão também são excluídos.

Destino de ingestão​

O ingestion_destination especifica o nome da tabela Delta de três partes onde os dados de transmissão são gravados.

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

Esquema de tabela de ingestão​

The ingestion table contains the message data along with metadata columns. The common columns are present for every source; the kafka_* columns are present only for a transmissão Kafka.

Coluna

Tipo

Origem

Descrição

key

Varia (de key_schema)

Comum

A key da mensagem, estruturada de acordo com o esquema que você forneceu.

value

Varia (de payload_schema)

Comum

O valor da mensagem (payload), estruturado de acordo com o esquema fornecido.

stream_record_timestamp

TIMESTAMP

Comum

The record Timestamp. For forward-fill data, this is the source ingest Timestamp. For backfill data, this is customer-supplied.

record_source

STRING

Comum

"stream" (preenchimento avançado a partir da transmissão ativa) ou "backfill" (a partir da origem de preenchimento retroativo).

kafka_topic

STRING

Kafka

O tópico do Kafka do qual o registro foi consumido.

kafka_partition

INT

Kafka

A partição do Kafka da qual o registro foi consumido.

kafka_offset

LONG

Kafka

O offset do Kafka do registro dentro de sua partição.

Coluna

Tipo

Origem

Descrição

key

Varia (de key_schema)

Comum

A key da mensagem, estruturada de acordo com o esquema que você forneceu.

value

Varia (de payload_schema)

Comum

O valor da mensagem (payload), estruturado de acordo com o esquema fornecido.

stream_record_timestamp

TIMESTAMP

Comum

The record Timestamp. For forward-fill data, this is the source ingest Timestamp. For backfill data, this is customer-supplied.

record_source

STRING

Comum

"stream" (preenchimento avançado a partir da transmissão ativa) ou "backfill" (a partir da origem de preenchimento retroativo).

kafka_topic

STRING

Kafka

O tópico do Kafka do qual o registro foi consumido.

kafka_partition

INT

Kafka

A partição do Kafka da qual o registro foi consumido.

kafka_offset

LONG

Kafka

O offset do Kafka do registro dentro de sua partição.

Fonte de preenchimento retroativo​

Como o pipeline de preenchimento avançado começa a partir da posição mais recente na origem, ele não captura mensagens que existiam antes da criação da transmissão. Para fornecer cobertura de dados históricos para o treinamento, configure uma origem de preenchimento retroativo opcional.

Quando uma fonte de preenchimento retroativo é configurada, o Databricks executa um Job único MERGE INTO que copia as linhas de preenchimento retroativo para a tabela de ingestão com record_source="backfill". A execução do MERGE só ocorre depois que o verificador de sobreposição confirma que a origem do preenchimento e a transmissão de preenchimento progressivo têm carimbos de data/hora sobrepostos (consulte Sobreposição entre o preenchimento e os dados da transmissão em tempo real). Se a condição de sobreposição não for atendida dentro de 2 dias, o merge realiza a execução de qualquer forma para evitar o bloqueio indefinido.

A tabela de preenchimento retroativo deve incluir uma coluna stream_record_timestamp do tipo TIMESTAMP no fuso horário UTC. Outras colunas de metadados são passadas se estiverem presentes na origem do preenchimento retroativo ou definidas como NULL, caso contrário. Para o Kafka, estes são kafka_topic, kafka_partition e kafka_offset.

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

Sobreposição entre dados de preenchimento retroativo e transmissão ao vivo​

Antes de executar um MERGE entre o preenchimento de dados e a tabela de ingestão, uma verificação de sobreposição compara os carimbos de data/hora nas duas tabelas:

  • **Máximo de preenchimento retroativo**: O máximo stream_record_timestamp na fonte de preenchimento retroativo.
  • Ingestão mín. : O stream_record_timestamp mínimo de linhas (record_source="stream") na tabela de ingestão.

O MERGE prossegue quando o carimbo de data/hora mais recente do preenchimento excede o carimbo de data/hora mais antigo da tabela de ingestão em pelo menos 1 hora. Essa sobreposição garante que não haja lacunas na tabela de ingestão. Se a condição de sobreposição não for atendida em 2 dias, a execução do merge ocorre de qualquer forma para evitar o bloqueio por tempo indeterminado.

Como o pipeline de ingestão começa a partir da posição mais recente na origem, ele captura apenas as mensagens que chegam após a criação da transmissão. Sua origem de preenchimento retroativo deve conter dados que se estendam até o intervalo de tempo de ingestão — não apenas até o momento de criação da transmissão.

Por exemplo, se você criar uma transmissão às 15:00, o pipeline de preenchimento progressivo começa a ler mensagens das 15:00 em diante. Sua fonte de preenchimento retroativo deve incluir dados com carimbos de data/hora até, pelo menos, às 16:00 (1 hora após o começar do preenchimento progressivo) para satisfazer a verificação de sobreposição. Isso significa que você deve atualizar sua tabela de preenchimento retroativo depois das 16:00 para garantir que a tabela de ingestão não tenha lacunas.

Eliminação de duplicação​

Use deduplication_columns para especificar caminhos de coluna para identificar linhas duplicadas durante a ingestão entre o preenchimento retroativo e o preenchimento progressivo de dados de transmissão. Use a notação de ponto para campos aninhados (por exemplo, "value.user_id").

Escolha as colunas de desduplicação com base em seus dados:

  • Se cada registro em sua transmissão contiver um identificador exclusivo (por exemplo, value.transaction_id), use essa coluna para deduplicação.
  • Se sua fonte de preenchimento retroativo incluir as colunas kafka_partition e kafka_offset, use-as para identificar exclusivamente cada registro.
  • Se nenhuma coluna de deduplicação for especificada, a key de deduplicação default será a combinação completa de key, value e stream_record_timestamp. Isso não é recomendado, pois a correspondência rigorosa de critérios pode facilmente levar a duplicatas.
Python
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)

Atribuição de custos​

Defina tags e budget_policy_id no IngestionConfig para atribuir o custo da ingestão gerenciada da transmissão. O Databricks os aplica ao pipeline Lakeflow de ingestão e aos seus jobs de preenchimento avançado (forward-fill) e preenchimento retroativo (backfill) quando o Stream é criado.

Para ver um exemplo, os limites de tags e como fazer query dos gastos atribuídos, consulte Atribuir custos com tags e políticas de uso serverless.

Excluir colunas de uma transmissão​

Use excluded_columns para remover colunas específicas de uma transmissão que você não deseja ingerir. Uma coluna excluída não é gravada na tabela de ingestão e não pode ser referenciada por um recurso ou usada no treinamento.

Specify each column using dot notation into the message key or value, such as value.user.email or key.account_id. These columns are dropped from the decoded key and value across ingestion, backfill, and materialization. If a path points to a struct, all of its nested fields are dropped as well (for example, value.address also drops value.address.city and value.address.zip).

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

When using direct schemas, the excluded column must already exist in the key or value schema, or create_stream fails. Ao usar um registro de esquema (schema registry), você pode excluir uma coluna antes que ela exista. Uma coluna excluída também não pode ser uma coluna de desduplicação, pois as colunas de desduplicação são obrigatórias para identificar linhas duplicadas. Qualquer recurso que faça referência a uma coluna excluída (por exemplo, como uma entidade, série temporal ou entrada) falha ao ser criado.

Você pode alterar as colunas excluídas de uma transmissão após a criação com update_stream, em transmissões de esquema direto e com suporte a registro de esquema. Consulte Update a transmissão para obter mais detalhes.

Gerenciar transmissões​

Obter uma transmissão​

Python
stream = client.get_stream(name="my_catalog.my_schema.my_stream")

Listar transmissões​

Python
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)

Defina include_schemas=True para incluir detalhes completos do esquema. Esquemas podem ser grandes e isso pode resultar em uma operação de longa duração. Para recuperar esquemas individualmente, use get_stream.

Atualizar uma transmissão​

Use update_stream para alterar uma transmissão após a criação. Passe schema_config para evoluir um esquema direto, excluded_columns para alterar quais colunas são descartadas ou ambos. A atualização de outros campos não é compatível. Em vez disso, crie uma nova transmissão.

A atualização de uma transmissão reinicia seu pipeline de ingestão para que a alteração entre em vigor. A ingestão geralmente é retomada em alguns minutos.

Evoluir um esquema direto​

Para um transmissão que usa esquemas diretos, passe DirectSchemas para schema_config. Defina payload_schema, key_schema ou ambos. Um lado que você não definir permanece inalterado. Transmissões com suporte a registro de esquemas rejeitam uma atualização de schema_config e devem ser evoluídos por meio do registro.

Python
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"}'
' }'
'}'
)
),
),
)

As atualizações de esquema devem ser compatíveis com versões anteriores para que o pipeline de ingestão em execução possa continuar decodificando os registros existentes e gravando na tabela de ingestão. Quaisquer outras alterações serão rejeitadas.

O que é permitido depende do formato:

  • JSON e Protobuf : adicionam campos opcionais, removem campos e ampliam o tipo de um campo (por exemplo, de int para bigint). O Protobuf também permite reordenar os campos.
  • Avro : permite apenas expandir int para long e remover um campo à direita cujos bytes nenhum campo posterior lê. Para evoluir um esquema Avro de forma mais livre, use uma transmissão apoiada por um registro de esquemas.

A adição de campos aumenta as estruturas key e value decodificadas da tabela de ingestão. As linhas gravadas antes da atualização mantêm o formato original, e os campos adicionados aparecem como NULL para essas linhas anteriores. A remoção e as alterações de tipo entram em vigor apenas para os registros ingeridos após a atualização.

Change excluded columns​

Passe o conjunto completo e novo de caminhos de coluna para excluded_columns, o que substitui o conjunto existente. Passe uma lista vazia ([]) para limpar todas as exclusões. Para obter detalhes sobre esse comportamento, consulte Excluir colunas de uma transmissão.

Python
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
excluded_columns=["value.user.email", "value.user.ssn"],
)

A alteração de colunas excluídas é exclusiva para o sentido de avanço. As colunas recém-excluídas param de ser gravadas (aparecendo em NULL) e as recém-incluídas começam a ser preenchidas daqui para frente, enquanto as linhas gravadas anteriormente permanecem inalteradas. Para evitar que uma nova coluna seja ingerida:

  • Schema registry : adicione a coluna a excluded_columns primeiro e aguarde o reinício do pipeline de ingestão; em seguida, registre a nova versão do esquema no registro.
  • Esquemas diretos : adicione a coluna a schema_config e a excluded_columns na mesma chamada de update_stream.

Excluir uma transmissão​

A exclusão de uma transmissão também exclui seu pipeline de ingestão e tabela de ingestão.

atenção

Quaisquer modelos ou recursos que referenciam a transmissão excluída não terão mais acesso aos dados da transmissão subjacente. Crie uma cópia da tabela de ingestão antes da exclusão se precisar desses dados, mas não precisar mais da transmissão.

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

Notebook de exemplo​

Para um exemplo completo que cria uma Transmissão, define recursos de transmissão e é implantado em um endpoint de serving, consulte o Notebook a seguir:

Notebook de início rápido de view de recurso de transmissão