Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Use a API de pipeline do Lakeflow sink com fluxos para gravar registros transformados por um pipeline em um destino de dados externo. Os destinos de dados externos incluem tabelas gerenciadas e externas do Unity Catalog e serviços de streaming de eventos, como Apache Kafka ou Hubs de Eventos do Azure. Você também pode usar coletores de dados para gravar em fontes de dados personalizadas escrevendo código Python para essa fonte de dados.
Para uma visão geral dos conceitos de sinks e de quando utilizá-los, consulte Sinks em pipelines do Lakeflow.
Observação
- A
sinkAPI só está disponível para Python. - Você pode criar um coletor personalizado com a API ForEachBatch. Consulte Usar ForEachBatch para gravar em coletores de dados arbitrários em pipelines.
Fluxo de trabalho do coletor
À medida que dados de eventos são ingeridos de uma fonte de streaming no seu pipeline, você processa e refina esses dados em transformações nesse pipeline. Em seguida, você usa o processamento de fluxos de acréscimo para transmitir os registros de dados transformados para um coletor. Você cria esse coletor usando a função create_sink(). Para obter mais detalhes sobre a função create_sink, consulte a referência da API do coletor.
Se você tiver um pipeline que crie ou processe seus dados de eventos de streaming e prepare registros de dados para gravação, estará pronto para usar um coletor.
A implementação de um coletor consiste em duas etapas:
- Crie o coletor.
- Use um fluxo de acréscimo ou fluxo de atualização para gravar os registros preparados no coletor.
Criar um coletor
O Databricks dá suporte a vários tipos de coletores de destino nos quais você grava seus registros processados de dados transmitidos:
- Coletores de tabelas Delta (incluindo tabelas externas e gerenciadas do Catálogo do Unity)
- Coletores do Apache Kafka
- Coletores dos Hubs de Eventos do Azure
- Coletores personalizados escritos em Python, utilizando fontes de dados personalizadas do Python
Abaixo estão exemplos de configurações para coletores do Delta, Kafka e Hubs de Eventos do Azure e fontes de dados personalizadas do Python:
Coletores Delta
Para criar um coletor Delta por caminho de arquivo:
dp.create_sink(
name = "delta_sink",
format = "delta",
options = {"path": "/Volumes/catalog_name/schema_name/volume_name/path/to/data"}
)
Para criar um coletor Delta por nome de tabela usando um catálogo e caminho de esquema totalmente qualificados:
dp.create_sink(
name = "delta_sink",
format = "delta",
options = { "tableName": "catalog_name.schema_name.table_name" }
)
Coletores do Kafka e Hubs de Eventos do Azure
Este código funciona tanto para sinks do Apache Kafka quanto do Hubs de Eventos do Azure.
credential_name = "<service-credential>"
eh_namespace_name = "dp-eventhub"
bootstrap_servers = f"{eh_namespace_name}.servicebus.windows.net:9093"
topic_name = "dp-sink"
dp.create_sink(
name = "eh_sink",
format = "kafka",
options = {
"databricks.serviceCredential": credential_name,
"kafka.bootstrap.servers": bootstrap_servers,
"topic": topic_name
}
)
O credential_name é uma referência a uma credencial de serviço do Catálogo do Unity. Para obter mais informações, consulte Usar as credenciais de serviço do Catálogo do Unity para se conectar a serviços de nuvem externos.
Fontes de dados personalizadas do Python
Supondo que você tenha uma fonte de dados personalizada do Python registrada como my_custom_datasource, o código a seguir pode gravar nessa fonte de dados.
from pyspark import pipelines as dp
# Assume `my_custom_datasource` is a custom Python streaming
# data source that writes data to your system.
# Create Lakeflow pipelines sink using my_custom_datasource
dp.create_sink(
name="custom_sink",
format="my_custom_datasource",
options={
<options-needed-for-custom-datasource>
}
)
# Create append flow to send data to RequestBin
@dp.append_flow(name="flow_to_custom_sink", target="custom_sink")
def flow_to_custom_sink():
return read_stream("my_source_data")
Para obter detalhes sobre como criar fontes de dados personalizadas no Python, consulte fontes de dados personalizadas do PySpark.
Para obter mais detalhes sobre como usar a função create_sink, confira a referência de API do coletor.
Depois que o coletor for criado, você poderá começar a transmitir registros processados para ele.
Gravar em um coletor com um fluxo de acréscimo
Com o coletor criado, a próxima etapa é gravar registros processados nele especificando-o como o destino para a saída de registros por um fluxo de acréscimo. Faça isso especificando o coletor como o valor target no decorador append_flow.
- Para tabelas gerenciadas e externas do Catálogo do Unity, use o formato
deltae especifique o caminho ou o nome da tabela nas opções. Seu pipeline deve ser configurado para usar o Catálogo do Unity. - Para tópicos do Apache Kafka, use o formato
kafkae especifique o nome do tópico, as informações de conexão e as informações de autenticação nas opções. Essas são as mesmas opções às quais um coletor do Kafka de Streaming Estruturado do Spark dá suporte. Confira Configurar o gravador de Streaming Estruturado do Kafka. - Para Hubs de Eventos do Azure, use o formato
kafkae especifique o nome dos Hubs de Eventos, as informações de conexão e as informações de autenticação nas opções. Essas são as mesmas opções com suporte em um coletor de Hubs de Eventos de Streaming Estruturado do Spark que usa a interface do Kafka. Consulte Autenticação.
Abaixo estão exemplos de como configurar fluxos para gravar em coletores do Delta, Kafka e Hubs de Eventos do Azure com registros processados pelo seu pipeline.
Coletor Delta
@dp.append_flow(name = "delta_sink_flow", target="delta_sink")
def delta_sink_flow():
return(
spark.readStream.table("spark_referrers")
.selectExpr("current_page_id", "referrer", "current_page_title", "click_count")
)
Coletores do Kafka e Hubs de Eventos do Azure
@dp.append_flow(name = "kafka_sink_flow", target = "eh_sink")
def kafka_sink_flow():
return (
spark.readStream.table("spark_referrers")
.selectExpr("cast(current_page_id as string) as key", "to_json(struct(referrer, current_page_title, click_count)) AS value")
)
O parâmetro value é obrigatório para um coletor dos Hubs de Eventos do Azure. Parâmetros adicionais como key, partition, headerse topic são opcionais.
Para obter mais detalhes sobre o decorador append_flow, consulte Fluxos padrão e fluxos anexados.
Limitações
Há suporte apenas para a API do Python. Não há suporte para SQL.
Há suporte apenas para consultas de streaming. Não há suporte para consultas em lote.
Somente
append_floweupdate_flowpodem ser usados para gravar em coletores. Outros fluxos, comocreate_auto_cdc_flow, não têm suporte e você não pode usar um coletor em uma definição de conjunto de dados de pipeline. Por exemplo, não há suporte para o seguinte:@table("from_sink_table") def fromSink(): return read_stream("my_sink")Para os coletores Delta, o nome da tabela deve ser totalmente qualificado. Especificamente, para tabelas externas gerenciadas do Unity Catalog, o nome da tabela deve ser no formato
<catalog>.<schema>.<table>. Para o metastore do Hive, ele deve estar no formato<schema>.<table>.A execução de uma atualização completa não limpa os dados de resultados computados anteriormente nos coletores. Isso significa que todos os dados reprocessados são acrescentados ao coletor e os dados existentes não são alterados.
Não há suporte para expectativas de pipelines.
O controle de saída sem servidor dá suporte apenas a conectores de coletores do Kafka e do Delta Lake. Veja O que é o controle de saída sem servidor?.