Configurar um fluxo

Importante

Esse recurso está em Visualização Pública. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.

Um Stream representa uma fonte de dados de streaming externa, como o Apache Kafka. Os fluxos armazenam detalhes de conexão, autenticação, esquemas e configuração de ingestão. Depois que um stream é criado, você pode fazer referência a ele em definições de Feature View para criar features de streaming em tempo real.

Os fluxos têm nomes de três partes (catalog.schema.stream_name). O acesso a um Stream é regido por sua tabela de ingestão associada. Veja Ingestão e preenchimento retroativo para mais detalhes.

Requirements

  • Para executar comandos de notebook: sem servidor ou um cluster de computação clássico executando o Databricks Runtime 17.0 ML ou superior.
  • O pacote Python feature-engineering-client versão 0.18.0 ou superior deve estar instalado.

Conectando-se a fontes de streaming

Antes de definir funcionalidades de streaming, conecte e teste uma conexão de pipeline de streaming do Lakeflow com seu broker do Kafka. O Feature Store usa SDP serverless, o que significa que você precisará de um mecanismo para conectar sua computação clássica (broker ou endpoint) à computação serverless do Databricks. Isso é feito por meio de produtos como privatelink ou permitindo que seu computador clássico seja acessível pela internet pública.

Criar um fluxo

Use create_stream() para criar um novo Stream. Um Stream requer quatro componentes de configuração:

  • Configuração de fonte: Especifica a plataforma de streaming e detalhes específicos da fonte, como a assinatura do tópico para uma fonte Kafka.
  • Configuração de conexão: especifica como se conectar e autenticar à plataforma de streaming, incluindo servidores de inicialização e credenciais.
  • Configuração de esquema: define a estrutura de chaves e valores de mensagem.
  • Configuração de ingestão: especifica onde e como os dados de fluxo são ingeridos. Veja Ingestão e preenchimento retroativo para mais detalhes.

Para obter a definição de conexão e source_config específicos, junto com um exemplo create_stream() completo, consulte Apache Kafka. O esquema e as opções de ingestão são compartilhados entre as fontes.

Apache Kafka

Para fazer streaming do Apache Kafka, use KafkaStreamConfig como a configuração de origem e uma conexão com o Unity Catalog para autenticação. Veja Streaming em computação serverless e Conecte-se ao Apache Kafka para conectividade 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"
        ),
    ),
)

Modos de assinatura Kafka

O modo de assinatura especifica como o Stream seleciona os tópicos do Kafka dos quais consumirá. Há suporte para três modos:

Modo Description Exemplo
subscribe Lista separada por vírgulas de nomes de tópico KafkaSubscriptionMode(subscribe="topic1,topic2")
subscribe_pattern Java: correspondência de padrões regex em nomes de tópicos KafkaSubscriptionMode(subscribe_pattern="events-.*")
assign JSON especificando atribuições de partição de tópico KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Autenticação Kafka

Use uma conexão do Unity Catalog para autenticar no cluster do Kafka. Essa é a abordagem recomendada para autenticação gerenciada. Para criar uma conexão, consulte Criar uma conexão. O criador do Stream deve ter USE CONNECTION na conexão. Qualquer usuário que materialize funcionalidades usando o Stream como origem também deve ter USE CONNECTION na conexão.

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

A conexão suporta tanto autenticação IAM (credencial de serviço) quanto autenticação SASL.

IAM (credencial de serviço)

Autentique com uma credencial de serviço do Unity Catalog, por exemplo, para 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 credential opção:

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 do Stream e a qualquer identidade que materialize recursos com o Stream. Veja Conceder permissões para usar uma credencial de serviço para acessar um serviço de nuvem externo.

SASL

A autenticação SASL utiliza 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 utiliza SASL/SCRAM. Para SASL/PLAIN, defina sasl_mechanism para 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 direto

Para autenticação direta do mTLS, forneça arquivos de repositório de chaves e de armazenamento confiável armazenados em um volume do Unity Catalog, com senhas referenciadas por meio de escopos secretos do Databricks. Para obter mais informações sobre a autenticação SSL com o Kafka, consulte Usar o SSL para se conectar Azure Databricks ao 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"
        ),
    ),
)

Configuração do esquema

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

Cada um SchemaConfig aceita um de três formatos, correspondendo à forma como a fonte serializa suas mensagens: json_schema, avro_schema, ou proto_schema. Se nenhum esquema for fornecido para uma chave ou conteúdo, ele será tratado como uma cadeia de caracteres 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 externo de esquemas, veja Registro de esquemas para detalhes.

Esquema JSON

Forneça uma string JSON Schema para 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"}'
    ),
)

Esquema Avro

Forneça uma string do esquema Avro para avro_schema. Tipos lógicos Avro são suportados, incluindo timestamp-millis, date, e 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"}}'
            '  ]'
            '}'
        )
    ),
)

Esquema Protobuf

Forneça um ProtoSchemaSpec a proto_schema com o texto-fonte de Buffers de Protocolo.proto e o nome da mensagem de payload. Importe ProtoSchemaSpec de databricks.feature_engineering.entities.

message_name deve ser o nome da mensagem totalmente qualificada, incluindo o package declarado no .proto texto (por exemplo, com.example.Event, não Event). Tanto a sintaxe proto2 quanto a proto3 são suportadas.

google.protobuf.Timestamp e os tipos de wrapper escalar (StringValue, Int32Value, e assim por diante) são suportados, e suas importações são resolvidas automaticamente. Outros tipos bem conhecidos, como Duration, Struct, e Any, são rejeitados; codifica esses valores como um escalar ou mensagem suportada. Os tipos escalares fixed32 e fixed64, assim como map com chaves que não são strings, também não são suportados.

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

A Databricks decodifica cada mensagem com as funções from_json, from_avro e from_protobuf do Spark. Os comportamentos a seguir se aplicam quer você declare o esquema diretamente, quer o resolva a partir de um registro de esquemas:

  • Registros malformados. A decodificação usa o PERMISSIVE modo, então um registro que não corresponde ao seu esquema decodifica para um valor nulo em vez de falhar o fluxo.
  • Uniões do Avro. Uma união de múltiplos tipos de registro decodifica para uma estrutura com um campo por tipo de registro, cada um nomeado em referência ao seu registro Avro.
  • Tipo Protobuf. Inteiros sem sinal são decodificados como um tipo assinado mais amplo (por exemplo, uint32 para BIGINT e uint64 para DECIMAL(20,0)), campos de enumeração são decodificados para o nome da string correspondente, e tipos encapsuladores escalares (por exemplo, StringValue e Int32Value) são decodificados como uma coluna anulável do tipo encapsulado.

Registro de esquema

Os registros de esquemas armazenam e versionam esquemas que produtores e consumidores de streaming utilizam, aplicando regras de compatibilidade conforme esses esquemas evoluem. Quando um registro de esquema externo é configurado, a Feature Store lê o esquema do registro e o usa para decodificar a mensagem de streaming. Você não declara o esquema inline no Stream ao usar um registro de esquemas.

O suporte ao registro de esquemas possui as seguintes limitações:

  • Compatível apenas com streams do Kafka.
  • Apenas o Registro de Esquema Confluente é suportado
  • Apenas os formatos Avro e Protobuf são suportados. Para ler mensagens JSON, declare o esquema diretamente, em vez disso. Veja o esquema JSON.
  • Cada Stream está conectado exatamente a um sujeito Confluente para o valor da mensagem e a um para a chave da mensagem (se fornecida). Tópicos de fluxo contendo múltiplos registros de esquema não são uma configuração suportada. Se seu Stream se conecta a tópicos que contêm múltiplos esquemas, registros que não correspondem ao esquema do assunto especificado são decodificados como nulos.

Conecte-se a um registro de esquemas

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 escopo secreto Databricks. A identidade de execução do Stream deve ter a permissão READ no escopo de segredo, pois o pipeline de ingestão lê o segredo em runtime. Para como criar e configurar uma conexão, veja Criar uma conexão.

Adicione as schema_registry_urlopções , 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 se autentica no broker com uma credencial de serviço do Unity Catalog e no registro com uma chave de 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>')
)

Defina a opção schema_registry_api_secret na conexão do Kafka e a referência ao escopo do segredo no Stream para o mesmo segredo.

Crie um fluxo que use um registro de esquema

Passe a SchemaRegistryConfig como o schema_config. Consulte 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 chave da mensagem. Pelo menos um localizador deve ser fornecido.

Note as diferenças aqui em comparação com os exemplos diretos de esquemas na seção de configuração de esquemas . Ao usar um registro de esquemas, você não fornece o esquema diretamente no stream para schema_config. Em vez disso, você especifica um SchemaRegistryConfig que identifica o esquema no registro.

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 sujeito Confluente é o escopo nomeado sob o qual o histórico de versões de um esquema é registrado e a compatibilidade é aplicada. Defina subject como o nome do escopo relevante, que é normalmente determinado a partir da estratégia de nome do assunto:

  • TopicNameStrategy (padrão, deriva o sujeito do nome do tópico): <topic>-value para o valor e <topic>-key para a chave. Por exemplo, o esquema de valores do tema transactions usa o sujeito transactions-value.
  • RecordNameStrategy (deriva o assunto do nome de registro do esquema, independentemente do tópico): o nome de 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 do tema e dos registros): <topic>-<fully-qualified-record-name>, como transactions-com.example.Payment.

format é obrigatório. Defina para SchemaLocatorFormat.FORMAT_AVRO ou SchemaLocatorFormat.FORMAT_PROTOBUF para combinar com a forma como o tema é serializado.

Evolução do esquema

O pipeline de ingestão resolve o esquema atual do sujeito quando ele começa. Quando você registra uma nova versão de esquema retrocompatível no assunto no registro de esquemas, o pipeline em execução continua usando a versão com a qual começou.

Para Streams suportados por registro de esquemas, o pipeline de ingestão reinicia automaticamente a cada poucas horas. A cada reinicialização, o Stream obtém a versão mais recente do esquema do tópico e os campos novos ou alterados aparecem na tabela de ingestão.

Os Streams que usam esquemas diretos em vez de um registro de esquemas evoluem seu esquema com update_stream. Consulte Atualizar um Stream.

Para saber como o pipeline lida com registros que não correspondem ao esquema que está usando atualmente, veja Decodificação de dados usando esquemas.

Filtrar registros por tipo

Um Stream decodifica cada registro usando um esquema de chave e valor único (se fornecido), seja especificado diretamente ou usando um registro de esquemas. Como um tópico pode conter mais de um tipo de registro e os Streams podem se inscrever em vários tópicos, use record_type_filter para selecionar quais registros do tópico pertencem a este Stream.

Forneça uma expressão SQL que referencie os campos decodificados com notação de ponto, por exemplo, value.event_type = 'transaction'. 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 um Stream para outros tipos de registro, crie um Stream separado com um record_type_filterdiferente.

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

Mesmo sem record_type_filter, a decodificação nunca falha no Stream. Um registro que não corresponde ao esquema configurado é decodificado permissivamente. Os registros são decodificados de uma das seguintes maneiras:

  • Em uma linha com NULL valores para os campos que o esquema espera, mas o registro omite (JSON, Avro e Protobuf).
  • Em uma linha contendo valores que pertencem a um tipo de registro diferente (apenas 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 Streams separados é recomendado quando os esquemas diferem substancialmente entre os tipos de registro no tópico ou quando você deseja controlar o acesso a cada tipo de registro de forma independente. Para manter os custos sob controle, o Databricks recomenda que você mantenha um número pequeno de Streams, já que cada Stream possui um pipeline de ingestão e uma tabela de ingestão separados. Cada Stream também utiliza computação separada no momento da materialização. Você pode usar filtros específicos de recursos para materialização.

record_type_filter é diferente de filter_condition de um recurso. record_type_filter é definido no Stream e controla quais registros são ingeridos e disponibilizados para todos os recursos que usam o Stream 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 streaming para obter mais detalhes sobre filter_condition.

Ingestão e preenchimento retroativo

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

O acesso a um Stream é regido pela tabela de ingestão:

  • SELECT na tabela de ingestão concede acesso de leitura ao Stream.
  • 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.

Canal de ingestão

Quando um fluxo é criado, o Databricks inicia um pipeline de ingestão gerenciado que lê continuamente mensagens do fluxo de origem e as grava em uma tabela Delta (a tabela de ingestão). O pipeline começa na posição mais recente na fonte e roda continuamente, capturando apenas novas mensagens que chegam após a criação do fluxo. Essa tabela de ingestão é usada para treinamento com recursos de streaming. Quando um fluxo é excluído, 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 em três partes na qual os dados de streaming são gravados.

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

Esquema de tabela de ingestão

A tabela de ingestão contém os dados da mensagem junto com as colunas de metadados. As colunas comuns estão presentes para todas as fontes; as colunas kafka_* estão presentes apenas para um stream do Kafka.

Coluna Tipo Source Description
key Varia (de key_schema) Common A chave da mensagem, estruturada de acordo com o esquema que você forneceu.
value Varia (de payload_schema) Common O valor da mensagem (payload), estruturado de acordo com o esquema que você forneceu.
stream_record_timestamp TIMESTAMP Common A data e hora do registro. Para dados com preenchimento forward-fill, este é o timestamp de ingestão da origem. No caso de dados de backfill, eles são fornecidos pelo cliente.
record_source STRING Common Ou "stream" (preenchimento progressivo direto do fluxo em tempo real) ou "backfill" (da fonte de provisionamento).
kafka_topic STRING Kafka O tópico Kafka do qual o registro foi consumido.
kafka_partition INT Kafka A partição Kafka da qual o registro foi consumido.
kafka_offset LONG Kafka O deslocamento Kafka do registro dentro de sua partição.

Origem do preenchimento retroativo

Como o pipeline de preenchimento progressivo começa na posição mais recente na origem,, ele não captura mensagens anteriores à criação da transmissão. Para fornecer cobertura de dados históricos para treinamento, configure uma fonte de backfill opcional.

Quando uma fonte de provisionamento é configurada, o Databricks executa um trabalho único MERGE INTO que copia linhas de provisionamento para a tabela de ingestão com record_source="backfill". O MERGE é executado somente depois que o verificador de sobreposição confirma que a origem do provisionamento e o fluxo de preenchimento avançado têm carimbos de data/hora sobrepostos (consulte Sobreposição entre os dados de provisionamento e de transmissão ao vivo). Se a condição de sobreposição não for atendida em até 2 dias, o MERGE é executado mesmo assim para evitar bloqueio indefinido.

A tabela de provisionamento deve incluir uma coluna stream_record_timestamp do tipo TIMESTAMP com fuso horário UTC. Outras colunas de metadados são passadas se estiverem presentes na fonte de preenchimento, ou configuradas para NULL outra forma. Para Kafka, são kafka_topic, kafka_partition, e 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"
    ),
)

Sobreposição entre preenchimento retroativo e dados em fluxo ao vivo

Antes de executar um MERGE entre o provisionamento e a tabela de ingestão, uma verificação de sobreposição compara os timestamps das duas tabelas:

  • Provisionamento máximo: o máximo stream_record_timestamp na origem do provisionamento.
  • Ingestão mínima: o número mínimo stream_record_timestamp de linhas (record_source="stream") na tabela de ingestão.

A operação MERGE é executada quando o timestamp mais recente do provisionamento excede o timestamp 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 até 2 dias, o MERGE é executado mesmo assim para evitar bloqueio indefinido.

Como o pipeline de ingestão começa na posição mais recente da origem, ele só captura mensagens que chegam após a criação do fluxo de dados. Sua fonte de provisionamento deve conter dados que abranjam o intervalo de tempo de ingestão, não apenas até o momento de criação do fluxo.

Por exemplo, se você criar um stream às 15h, o pipeline de preenchimento progressivo começará a ler mensagens a partir das 15h. Sua fonte de provisionamento deve incluir dados com registros de data e hora até pelo menos 16h (1 hora após o início do preenchimento progressivo) para atender à verificação de sobreposição. Isso significa que você deve atualizar sua tabela de backfill após as 16h para garantir que a tabela de ingestão não tenha lacunas.

Eliminação de duplicação

Use deduplication_columns para especificar os caminhos das colunas para identificar linhas duplicadas durante a ingestão entre dados de provisionamento e dados de fluxo de preenchimento progressivo. Usar notação de ponto para campos aninhados (por exemplo, "value.user_id").

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

  • Se cada registro em seu fluxo contiver um identificador exclusivo (por exemplo, value.transaction_id), use essa coluna para eliminação de duplicação.
  • Se a fonte de provisionamento incluir as colunas kafka_partition e kafka_offset, use-as para identificar cada registro de forma exclusiva.
  • Se nenhuma coluna de eliminação de duplicação for especificada, a chave de eliminação de duplicação padrão será a combinação completa de key, valuee stream_record_timestamp. Isso não é recomendado, pois essa correspondência de critérios estritos pode facilmente levar a duplicatas.
ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
    deduplication_columns=["value.transaction_id"],
)

Atribuição de custo

Defina tags e budget_policy_id no IngestionConfig para atribuir os custos da ingestão gerenciada do Stream. O Azure Databricks aplica essas configurações ao pipeline de ingestão do Lakeflow e às suas tarefas de preenchimento (forward-fill e backfill) quando o Stream é criado.

Para ver um exemplo, os limites de tags e como consultar os gastos atribuídos, veja Atribuir custos com tags e políticas de uso sem servidor.

Excluir colunas de um Stream

Use excluded_columns para remover colunas específicas de um Stream 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.

Especifique cada coluna usando a notação de ponto na chave ou no valor da mensagem, como value.user.email ou key.account_id. Essas colunas são removidas dos valores decodificados key e value durante a ingestão, o provisionamento retroativo e a materialização. Se um caminho apontar para uma estrutura, todos os seus campos aninhados também serão removidos (por exemplo, value.address também remove value.address.city e value.address.zip).

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

Ao usar esquemas diretos, a coluna excluída já deve existir no esquema de chave ou valor, caso contrário, create_stream falha. Ao usar um registro de esquema, você pode excluir uma coluna antes que ela exista. Uma coluna excluída também não pode ser uma coluna de deduplicação, pois colunas de deduplicação são necessá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) não será criado.

Você pode alterar as colunas excluídas de um Stream após a criação com update_stream, tanto em Streams com esquema direto quanto em Streams com suporte de registro de esquema. Consulte Atualizar um Stream para obter mais detalhes.

Gerenciar fluxos

Obter um fluxo

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

Listar fluxos

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, o que pode resultar em uma operação demorada. Para recuperar esquemas individualmente, use get_stream.

Atualizar um Stream

Use update_stream para alterar um Stream 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. Crie um novo Stream.

A atualização de um Stream reinicia seu pipeline de ingestão, de modo que a alteração entre em vigor. A ingestão de dados normalmente é retomada dentro de alguns minutos.

Evoluir um esquema direto

Para um Stream que usa esquemas diretos, passe um DirectSchemas para schema_config. Defina payload_schema, key_schema, ou ambos. Um lado que você não definir permanecerá inalterado. Streams com suporte de registro de esquema rejeitam uma atualização schema_config e devem ser evoluídos por meio do registro.

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 retrocompatíveis 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 mudanças são rejeitadas.

O que é permitido depende do formato:

  • JSON e Protobuf: adicione campos opcionais, remova campos e amplie o tipo de um campo (por exemplo, int para bigint). O Protobuf também permite reordenar campos.
  • Avro: permite apenas o alargamento de int para long e a remoção de um campo final cujos bytes nenhum campo posterior lê. Para evoluir um esquema Avro com mais liberdade, use um fluxo com suporte de registro de esquema.

A adição de campos aumenta as estruturas decodificadas key e value da tabela de ingestão. As linhas escritas antes da atualização mantêm seu formato original e os campos adicionados são lidos como NULL para essas linhas anteriores. Remoções e mudanças de tipo só têm efeito para registros ingeridos após a atualização.

Alterar colunas excluídas

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

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 só pode ser aplicada a partir de agora. As colunas recém-excluídas param de ser escritas (aparecendo em NULL) e as recém-incluídas começam a ser preenchidas a partir de agora, enquanto as linhas escritas anteriormente permanecem como estão. Para impedir que uma nova coluna seja ingerida:

  • Registro de esquema: adicione a coluna a excluded_columns primeiro e aguarde a reinicialização 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 um fluxo

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

Warning

Todos os modelos ou recursos que fazem referência ao fluxo excluído não terão mais acesso aos dados de fluxo subjacentes. Crie uma cópia da tabela de ingestão antes da exclusão se você precisar desses dados, mas não precisar mais do fluxo.

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

Notebook de exemplo

Para ver um exemplo completo que cria um Stream, define recursos de streaming e faz a implantação em um endpoint de serviço, consulte o seguinte notebook:

Bloco de anotações de início rápido exibições de recursos de transmissão

Obter laptop