Configurar uma transmissão
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.
Um Stream representa uma fonte de dados de transmissão externa, como o Apache Kafka ou o Amazon Kinesis. As transmissões armazenam detalhes de conexão, autenticação, esquemas e configuração de ingestão. Depois que uma transmissão é criada, você pode referenciá-la usando definições de recurso 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 Python
feature-engineering-clientna versão 0.17.0 ou acima deve estar instalado.
Conectando-se a fontes de transmissão
Antes de definir recursos de transmissão, conecte e teste uma conexão de Lakeflow pipeline de transmissão ao seu broker do Kafka ou ao Endpoint de serviço do Kinesis. O 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 da Databricks. Isso é feito por meio de produtos como privatelink ou permitindo que seu compute clássico seja acessível a partir da Internet pública.
Para transmissão gerenciada pela AWS (Amazon MSK), consulte Conectividade privada Serverless ao Amazon MSK. Um padrão semelhante será necessário para outras clouds e para o Kinesis.
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.
Para obter o source_config específico da fonte e a configuração de conexão, juntamente com um exemplo create_stream() completo, consulte Apache Kafka ou Amazon Kinesis. As opções de esquema e ingestão são compartilhadas entre as fontes.
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.
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 |
|---|---|---|
| Lista de nomes de tópicos separados por vírgulas. |
|
| Padrão Java regex correspondente aos nomes de tópico |
|
| JSON especificando as atribuições de tópicos-partições |
|
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.
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:
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:
PLAINSCRAM-SHA-256SCRAM-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.
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.
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"
),
),
)
Amazon Kinesis
Para transmitir a partir de uma transmissão de dados do Amazon Kinesis, use KinesisStreamConfig como a configuração da fonte e uma conexão do Unity Catalog do tipo KINESIS para autenticação.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KinesisStreamConfig,
StreamNameList,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KinesisStreamConfig(
stream_names=StreamNameList(names=["my-kinesis-stream"]),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kinesis-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "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"
),
),
)
Identificadores de transmissão do Kinesis
Identifique a(s) transmissão(ões) de dados do Kinesis para ler com exatamente uma das seguintes opções. Uma única Transmissão pode ler de mais de uma transmissão de dados do Kinesis. Passe opções de origem adicionais por meio de extra_options — por exemplo, maxFetchRate para limitar a taxa de leitura por fragmento ou consumerMode="efo" para ler com o fan-out aprimorado (EFO) em vez do consumidor de polling default.
campo | Descrição | Exemplo |
|---|---|---|
| Lista de nomes de transmissões do Kinesis |
|
| Lista de ARNs de transmissão do Kinesis |
|
Autenticação do Kinesis
O Kinesis autentica-se por meio de uma conexão do Unity Catalog do tipo KINESIS que faz referência a uma credencial de serviço do Unity Catalog (um IAM role que concede acesso de leitura ao Kinesis) e à região da AWS do stream. Para criar uma conexão, consulte Autenticar com uma conexão do Unity Catalog. O criador da transmissão deve ter USE CONNECTION na conexão, assim como qualquer usuário que materialize recursos usando a transmissão como origem.
CREATE CONNECTION IF NOT EXISTS `my-kinesis-connection`
TYPE KINESIS
OPTIONS (
aws_region '<region>',
credential '<service_credential>'
)
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.
Para origens do Kinesis, forneça apenas payload_schema para os dados do registro. A chave de partição de uma mensagem do Kinesis é uma string de roteamento sem esquema, portanto key_schema não se aplica; a chave de partição ainda é capturada na coluna key da tabela de ingestão como uma string simples.
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.
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.
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.
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,
uint32paraBIGINTeuint64paraDECIMAL(20,0)), campos enum são decodificados para seu nome de string e tipos de wrapper escalar (por exemplo,StringValueeInt32Value) 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:
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.
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>-valuepara o valor e<topic>-keypara a chave. Por exemplo, o esquema de valor para o tópicotransactionsusa o assuntotransactions-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>, comotransactions-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.
Como o Databricks gerencia o pipeline de ingestão como um Lakeflow pipeline Serverless, o pipeline é reiniciado periodicamente. Na próxima reinicialização, ele captura a nova versão do esquema. Pode levar até uma semana para que campos novos ou alterados apareçam na tabela de ingestã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.
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:
SELECTna tabela de ingestão concede acesso de leitura à Transmissão.MANAGEna 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.
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 juntamente com as colunas de metadados. As colunas comuns estão presentes em todas as fontes; as colunas kafka_* estão presentes apenas para uma transmissão do Kafka, e as colunas kinesis_* apenas para uma transmissão do Kinesis.
Coluna | Tipo | Origem | Descrição |
|---|---|---|---|
| Varia (de | Comum | A key da mensagem, estruturada de acordo com o esquema que você forneceu. Uma transmissão do Kinesis não tem nenhum esquema de key, portanto, esta é a key de partição do registro como um |
| Varia (de | Comum | O valor da mensagem (payload), estruturado de acordo com o esquema fornecido. |
|
| Comum | The record Timestamp. For forward-fill data, this is the source ingest Timestamp. For backfill data, this is customer-supplied. |
|
| Comum |
|
|
| Kafka | O tópico do Kafka do qual o registro foi consumido. |
|
| Kafka | A partição do Kafka da qual o registro foi consumido. |
|
| Kafka | O offset do Kafka do registro dentro de sua partição. |
|
| Kinesis | A transmissão de dados do Kinesis da qual o registro foi consumido. |
|
| Kinesis | O fragmento do qual o registro foi consumido. |
|
| Kinesis | O número de sequência do registro dentro do seu fragmento. |
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.
Para o Kinesis, as colunas de metadados de passagem são kinesis_stream, kinesis_shard_id e kinesis_sequence_number.
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_timestampna fonte de preenchimento retroativo. - Ingestão mín. : O
stream_record_timestampmí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_partitionekafka_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,valueestream_record_timestamp. Isso não é recomendado, pois a correspondência rigorosa de critérios pode facilmente levar a duplicatas.
Para uma transmissão do Kinesis, use kinesis_shard_id e kinesis_sequence_number juntos para desduplicação — um número de sequência é exclusivo apenas dentro de um fragmento — além de kinesis_stream quando a Transmissão lê mais de uma transmissão do Kinesis.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Gerenciar transmissões
Obter uma transmissão
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Listar transmissões
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.
Excluir uma transmissão
A exclusão de uma transmissão também exclui seu pipeline de ingestão e tabela de ingestã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.
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: