Configurar uma transmissão

Importante

Este recurso está no Public Preview. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Visualizações . Ver Gerir as pré-visualizações de Azure Databricks.

Um Stream representa uma fonte externa de dados de streaming, como o Apache Kafka. Os fluxos armazenam detalhes de ligação, autenticação, esquemas e configuração de ingestão. Depois de um fluxo ser criado, pode referenciá-lo utilizando definições de Feature View para criar funcionalidades de streaming em tempo real.

Os cursos de água têm nomes em três partes (catalog.schema.stream_name). O acesso a um Stream é regulado pela sua tabela de ingestão associada. Consulte Ingestão e retropreenchimento para obter mais informações.

Requisitos

  • Para executar comandos de bloco de notas: sem servidor ou um cluster de computação clássico com o Databricks Runtime 17.0 ML ou superior.
  • O feature-engineering-client pacote Python versão 0.17.0 ou superior deve ser instalado.

Ligação às fontes de ribeiro

Antes de definir funcionalidades de streaming, estabeleça e teste uma ligação de pipeline Lakeflow de streaming ao seu broker Kafka. A Feature Store baseia-se em SDP serverless, o que significa que vai precisar de um mecanismo para ligar a sua computação clássica (broker ou endpoint) à computação serverless da Databrick. Isto é feito através de produtos como o privatelink ou permitindo que o seu computador clássico seja acessível a partir da internet pública.

Criar um fluxo

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

  • Configuração da fonte: Especifica a plataforma de streaming e detalhes específicos da fonte, como a subscrição do tópico para uma fonte Kafka.
  • Configuração de ligação: Especifica como se ligar e autenticar à plataforma de streaming, incluindo servidores bootstrap e credenciais.
  • Configuração do esquema: Define a estrutura das chaves e valores das mensagens.
  • Configuração de ingestão: Especifica onde e como os dados do fluxo são ingeridos. Consulte Ingestão e preenchimento retroativo para mais detalhes.

Para a configuração específica source_config da fonte e da ligação, juntamente com um exemplo completo create_stream() , veja Apache Kafka. O esquema e as opções de ingestão são partilhados entre as fontes.

Apache Kafka

Para transmitir a partir do Apache Kafka, use KafkaStreamConfig como configuração de origem e uma ligação ao Unity Catalog para autenticação. Consulte Streaming em computação sem servidor e Ligar-se ao Apache Kafka para 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 subscrição Kafka

O modo de subscrição especifica como o Stream seleciona os tópicos Kafka para consumir. São suportados três modos:

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

Autenticação Kafka

Use uma ligação ao Unity Catalog para autenticar no seu cluster Kafka. Esta é a abordagem recomendada para autenticação gerida. Para criar uma ligação, veja Criar uma ligação. O criador da transmissão deve ter USE CONNECTION na ligação. Qualquer utilizador que materialize elementos tendo o Stream como origem também tem de ter USE CONNECTION na ligação.

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

A ligação suporta tanto autenticação IAM (credencial de serviço) como autenticação SASL.

IAM (credencial de serviço)

Autentique-se com uma credencial de serviço do Unity Catalog, por exemplo para ligar à Amazon MSK com o 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 ligação, as identidades que utilizam a credencial de serviço precisam de ACCESS nela. Conceda ACCESS a credencial de serviço referenciada ao criador do Stream e a qualquer identidade que materialize funcionalidades com o Stream. Consulte 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 utilizador e palavra-passe. 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 ligação armazena estas credenciais de forma segura.

O exemplo seguinte 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 mTLS direta, forneça ficheiros keystore e truststore armazenados num volume do Unity Catalog, com palavras-passe referenciadas através dos escopos secretos do Databrick. Para mais informações sobre autenticação SSL com Kafka, veja Usar SSL para ligar o 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 valores das mensagens para que as definições de ingestão e funcionalidades possam ler campos individuais. Para fontes de Kafka, payload_schema corresponde ao valor da mensagem de Kafka (o value no modelo chave-valor de Kafka) e key_schema corresponde à chave de mensagem de Kafka. Pelo menos um dos payload_schema ou key_schema deve ser fornecido.

Cada um SchemaConfig aceita um de três formatos, correspondendo à forma como a fonte serializa as suas mensagens: json_schema, avro_schema, ou proto_schema. Se não for fornecido um esquema para uma chave ou carga útil, é tratado como uma cadeia simples.

Os exemplos de código nesta secção usam esquemas declarados em linha com DirectSchemas, onde o esquema é fornecido como uma cadeia. Para gerir esquemas usando um registo externo de esquemas, consulte Registo de esquemas para mais detalhes.

Esquema JSON

Forneça uma cadeia JSON Schema a 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 cadeia de caracteres de um esquema Avro a avro_schema. São suportados tipos lógicos avro, 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 de origem Protocol Buffers.proto e o nome da mensagem de carga útil. Importar 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 como a proto3 são suportadas.

google.protobuf.Timestamp e os tipos de wrapper escalar (StringValue, Int32Value, e assim sucessivamente) são suportados, e as suas importações são resolvidas automaticamente. Outros tipos bem conhecidos, como Duration, Struct, e Any, são rejeitados; codifica esses valores como escalar ou mensagem suportada em vez disso. Os tipos escalares fixed32 e fixed64 e map com chaves que não sejam cadeias de caracteres 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

O Databricks descodifica cada mensagem com as funções from_json, from_avro e from_protobuf do Spark. Os seguintes comportamentos aplicam-se quer declare o esquema em linha, quer o resolva a partir de um registo de esquemas:

  • Registos malformados. A decodificação usa o PERMISSIVE modo, por isso um registo que não corresponde ao seu esquema decodifica-se para um valor nulo em vez de falhar o fluxo.
  • Sindicatos Avro. Uma união de vários tipos de registo é descodificada numa estrutura com um campo por cada tipo de registo, cada campo com o nome do respetivo registo Avro.
  • Tipos Protobuf. Inteiros sem sinal descodificam para um tipo com sinal mais amplo (por exemplo, uint32 para BIGINT e uint64 para DECIMAL(20,0)), os campos de enumeração descodificam para o respetivo nome em cadeia de caracteres e os tipos de invólucro escalar (por exemplo, StringValue e Int32Value) descodificam para uma coluna que admite valores nulos do tipo encapsulado.

Registro de esquema

Os registos de esquemas armazenam e versionam esquemas que produtores e consumidores de streaming utilizam, aplicando regras de compatibilidade à medida que esses esquemas evoluem. Quando um registo de esquema externo é configurado, a Feature Store lê o esquema do registo e usa-o para decodificar a mensagem em streaming. Não declare o esquema em linha no Stream ao usar um registo de esquemas.

O suporte ao registo de esquemas tem as seguintes limitações:

  • Suportado apenas para os fluxos do Kafka.
  • Apenas o Confluent Schema Registry é suportado
  • Apenas os formatos Avro e Protobuf são suportados. Para ler mensagens JSON, declare o esquema em linha em vez disso. Ver esquema JSON.
  • Cada Stream está ligado a exatamente um sujeito Confluente para o valor da mensagem e um para a chave da mensagem (se fornecida). Tópicos de fluxo contendo múltiplos registos de esquema não são uma configuração suportada. Se o seu Stream se ligar a tópicos que contenham múltiplos esquemas, os registos que não correspondem ao esquema do assunto especificado são descodificados como nulos.

Liga-te a um registo de esquemas

Forneça os detalhes da ligação ao registo como opções na ligação Kafka Unity Catalog e armazene o segredo da API do registo num âmbito secreto Databricks. A identidade de execução do Stream tem de ter a permissão READ no escopo do segredo, porque o pipeline de ingestão lê o segredo em tempo de execução. Para saber como criar e configurar uma ligação, veja Criar uma ligação.

Adicione as opções schema_registry_url, schema_registry_api_key e schema_registry_api_secret à conexão usada na autenticação. O exemplo seguinte cria uma ligação Kafka que autentica o corretor com uma credencial de serviço do Unity Catalog e o registo com uma chave 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 tanto a opção schema_registry_api_secret na ligação ao Kafka como a referência ao âmbito do segredo no Stream para o mesmo segredo.

Crie um fluxo que utilize um registo de esquemas

Passe um(a) SchemaRegistryConfig como schema_config. Referenciar o segredo da API do registo com api_secret_ref, e identificar o assunto e o formato com payload_schema_locator para o valor da mensagem, ou key_schema_locator para a chave da mensagem. Deve ser fornecido pelo menos um localizador.

Note as diferenças aqui em comparação com os exemplos diretos de esquemas na secção de configuração de esquemas . Ao utilizar um registo de esquemas, não fornece o esquema diretamente no Stream para schema_config. Em vez disso, especifica um SchemaRegistryConfig que identifica o esquema no registo.

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

  • TopicNameStrategy (por defeito, deriva o assunto 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 utiliza o sujeito transactions-value.
  • RecordNameStrategy (deriva o assunto do nome de registo do esquema, independentemente do tópico): o nome de registo totalmente qualificado, como com.example.Payment. Este é o espaço de nomes e o nome do registo para Avro, ou o pacote e o nome da mensagem para Protobuf.
  • TopicRecordNameStrategy (combina os nomes do tema e do registo): <topic>-<fully-qualified-record-name>, como transactions-com.example.Payment.

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

Evolução do esquema

O pipeline de ingestão determina o esquema atual do assunto quando é iniciado. Quando regista uma nova versão de esquema compatível com versões anteriores no tema no registo do esquema, o pipeline em execução continua a usar a versão com que começou.

Como a Databricks gere o pipeline de ingestão como um pipeline Lakeflow sem servidor, este pipeline reinicia periodicamente. No reinício seguinte, detecta a nova versão do esquema. Pode demorar até uma semana para aparecerem campos novos ou alterados na tabela de ingestão.

Para saber como o pipeline lida com registos que não correspondem ao esquema que está a usar atualmente, veja Descodificar dados usando esquemas.

Ingestão e preenchimento retroativo

O ingestion_config parâmetro configura como os dados do fluxo são capturados e armazenados para treino 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 permissões de eliminação.

Para mais informações sobre privilégios da 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 gerido que lê continuamente mensagens do fluxo de origem e as escreve numa tabela Delta (a tabela de ingestão). O pipeline começa a partir da posição mais recente na origem e é executado continuamente, captando apenas as novas mensagens que chegam após a criação do fluxo. Esta tabela de ingestão é usada para treino com funcionalidades de streaming. Quando um fluxo é eliminado, o seu canal de ingestão e a tabela de ingestão também são eliminados.

Destino de ingestão

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

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

Esquema da 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 kafka_* colunas estão presentes apenas num fluxo Kafka.

Column Tipo Source Description
key Varia (a partir de key_schema) Common A chave da mensagem, estruturada de acordo com o esquema que forneceste.
value Varia (a partir de payload_schema) Common O valor da mensagem (payload), estruturado de acordo com o esquema que forneceste.
stream_record_timestamp TIMESTAMP Common O carimbo temporal do registo. Para dados de preenchimento direto, este é o carimbo temporal de ingestão de origem. Para dados retroativos, estes são fornecidos pelo cliente.
record_source STRING Common Ou "stream" (preenchimento progressivo a partir do fluxo em direto) ou "backfill" (a partir da fonte de preenchimento retroativo).
kafka_topic STRING Kafka O tema Kafka do qual o disco foi consumido.
kafka_partition INT Kafka A partição Kafka de onde o disco foi consumido.
kafka_offset LONG Kafka O deslocamento de Kafka do disco dentro da sua partição.

Fonte de preenchimento

Como o pipeline de forward-fill começa na posição mais recente na origem, não capta mensagens que já existiam antes de o stream ser criado. Para fornecer cobertura histórica de dados para treino, configure uma fonte opcional de preenchimento.

Quando é configurada uma origem de backfill, o Databricks executa uma tarefa única MERGE INTO que copia as linhas de backfill para a tabela de ingestão com record_source="backfill". O MERGE só funciona depois de o verificador de sobreposição confirmar que a fonte de preenchimento e o fluxo de preenchimento direto têm carimbos temporais sobrepostos (ver Sobreposição entre dados de preenchimento e transmissão em direto). Se a condição de sobreposição não for satisfeita no prazo de 2 dias, a MERGE é executada ainda assim para evitar bloqueios indefinidos.

A tabela de preenchimento deve incluir uma stream_record_timestamp coluna do tipo TIMESTAMP no fuso horário UTC. Outras colunas de metadados são transmitidas, se estiverem presentes na origem do preenchimento retroativo, ou definidas como NULL caso contrário. Para Kafka, estes 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 os dados de preenchimento e transmissão ao vivo

Antes de executar uma operação MERGE entre o backfill e a tabela de ingestão, uma verificação de sobreposição compara as marcas temporais das duas tabelas:

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

A operação MERGE é executada quando a marca temporal mais recente do preenchimento retroativo é superior à marca temporal mais antiga da tabela de ingestão em, pelo menos, 1 hora. Esta sobreposição garante que não existem lacunas na tabela de ingestão. Se a condição de sobreposição não for satisfeita no prazo de 2 dias, a MERGE é executada ainda assim para evitar bloqueios indefinidos.

Como o pipeline de ingestão começa na posição mais recente na fonte, só capta mensagens que chegam após a criação do fluxo. A sua origem de preenchimento retroativo deve conter dados que abranjam o intervalo de tempo de ingestão — e não apenas até ao momento de criação do fluxo.

Por exemplo, se criar um fluxo de dados às 15:00, a canalização de preenchimento progressivo começa a ler mensagens a partir das 15:00. A sua origem do preenchimento retroativo deve incluir dados com marcas temporais até, pelo menos, às 16:00 (1 hora após o início do preenchimento progressivo) para cumprir a verificação de sobreposição. Isto significa que deve atualizar a sua tabela de preenchimento após as 16:00 para garantir que a tabela de ingestão não tem lacunas.

Deduplication

Utilize deduplication_columns para especificar os caminhos das colunas para identificar linhas duplicadas durante a ingestão entre dados de fluxo de backfill e de forward-fill. Use notação de pontos para campos aninhados (por exemplo, "value.user_id").

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

  • Se cada registo no seu fluxo contiver um identificador único (por exemplo, value.transaction_id), use essa coluna para a deduplicação.
  • Se a sua origem de preenchimento retroativo incluir as colunas kafka_partition e kafka_offset, utilize-as para identificar cada registo de forma única.
  • Se não forem especificadas colunas de deduplicação, a chave de deduplicação por defeito é a combinação completa de key, value, e stream_record_timestamp. Isto não é recomendado, pois esta correspondência rigorosa de critérios pode facilmente levar a duplicados.
ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
    deduplication_columns=["value.transaction_id"],
)

Gerir fluxos

Obter uma transmissão

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 os detalhes completos do esquema. Os esquemas podem ser grandes e isso pode resultar numa operação de longa duração. Para recuperar esquemas individualmente, use get_stream.

Excluir uma transmissão

Eliminar um fluxo também elimina o seu pipeline de ingestão e a tabela de ingestão.

Warning

Quaisquer modelos ou funcionalidades que referenciam o fluxo eliminado deixarão de ter acesso aos dados subjacentes do fluxo. Crie uma cópia da tabela de ingestão antes de eliminar se precisar destes dados mas já não precisar do fluxo.

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

Bloco de notas de exemplo

Para um exemplo de ponta a ponta que cria um Stream, define funcionalidades de streaming e é implementado num endpoint de serviço, veja o seguinte caderno:

Notebook de introdução rápida do Streaming Feature Views

Obter caderno