使用 Lakeflow 管道 sink API 和 流,将经管道转换后的记录写入外部数据接收端。 外部数据接收器包括 Unity 目录托管表和外部表,以及 Apache Kafka 或 Azure 事件中心等事件流服务。 还可以通过编写 Python 代码,将数据接收器用于写入自定义数据源。
有关接收器概念的概述以及何时使用它们,请参阅 Lakeflow 管道中的接收器。
注释
- API
sink仅适用于 Python。 - 可以使用 ForEachBatch API 创建自定义接收器。 请参阅 使用 ForEachBatch 写入管道中的任意数据接收器。
接收端工作流
当事件数据从流式数据源引入管道时,可以通过管道转换处理和优化该数据。 然后使用追加流处理过程将转换后的数据记录传输到接收器。 使用 create_sink() 函数创建此接收器。 有关 create_sink 函数的更多详细信息,请参阅 Sink API 参考。
如果您的管道创建或处理流式事件数据,并准备好数据记录供写入到接收端,那么您就可以使用数据接收端。
实现接收器包括两个主要步骤:
创建下沉器
Databricks 支持将从流数据处理的记录写入几种不同类型的目标存储:
- Delta 表接收器(包括 Unity Catalog 托管表和外部表)
- Apache Kafka 接收器
- Azure 事件中心接收器
- 使用 Python 编写的自定义接收器,基于 Python 自定义数据源
下面是 Delta、Kafka 和 Azure 事件中心接收器和 Python 自定义数据源的配置示例:
德尔塔水槽
若要按文件路径创建 Delta 接收器,请执行以下操作:
dp.create_sink(
name = "delta_sink",
format = "delta",
options = {"path": "/Volumes/catalog_name/schema_name/volume_name/path/to/data"}
)
若要使用完全限定的目录和架构路径按表名称创建 Delta 接收器,请执行以下操作:
dp.create_sink(
name = "delta_sink",
format = "delta",
options = { "tableName": "catalog_name.schema_name.table_name" }
)
Kafka 和 Azure 事件中心接收器
此代码适用于 Apache Kafka 和 Azure 事件中心接收器。
credential_name = "<service-credential>"
eh_namespace_name = "dp-eventhub"
bootstrap_servers = f"{eh_namespace_name}.servicebus.chinacloudapi.cn: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
}
)
credential_name 是对 Unity Catalog 服务凭据的引用。 有关详细信息,请参阅 使用 Unity 目录服务凭据连接到外部云服务。
Python 自定义数据源
假设你注册了 my_custom_datasourcePython 自定义数据源,则以下代码可以写入该数据源。
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")
有关在 Python 中创建自定义数据源的详细信息,请参阅 PySpark 自定义数据源。
有关使用 create_sink 函数的更多详细信息,请参阅 接收器 API 参考。
创建接收器后,可以开始将处理过的记录流式传输到接收器。
使用追加流写入接收器
创建接收器后,下一步是通过将接收器指定为追加流输出的记录的目标,向其写入处理后的记录。 为此,请将接收器指定为 target 修饰器中的 append_flow 值。
- 对于 Unity 目录托管表和外部表,请使用格式
delta并在选项中指定路径或表名称。 您的数据管道需要配置以使用 Unity Catalog。 - 对于 Apache Kafka 主题,请使用
kafka格式,并在选项中指定主题名称、连接信息和身份验证信息。 这些选项与 Spark 结构化流式处理 Kafka 接收器支持的选项相同。 请参阅配置 Kafka 结构化流式处理写入器。 - 对于 Azure 事件中心,请使用
kafka格式,并在选项中指定事件中心名称、连接信息和身份验证信息。 这些选项与使用 Kafka 接口的 Spark 结构化流式处理事件中心接收器中支持的选项相同。 请参阅身份验证。
下面是一些示例,介绍如何设置流以将记录写入 Delta、Kafka 和 Azure 事件中心接收器,并由您的管道处理这些记录。
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")
)
Kafka 和 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")
)
对于 Azure 事件中心接收器,value 参数是必需的。 其他参数(如 key、partition、headers和 topic)是可选的。
有关 append_flow 装饰器的更多信息,请参阅 默认流和追加流。
局限性
仅支持 Python API。 不支持 SQL。
仅支持流式处理查询。 不支持批处理查询。
只有
append_flow和update_flow可用于向接收器写入。 不支持其他流,例如create_auto_cdc_flow,且不能在管道数据集定义中使用汇集器。 例如,不支持以下各项:@table("from_sink_table") def fromSink(): return read_stream("my_sink")对于 Delta 接收器,表名称必须完全限定。 具体而言,对于 Unity Catalog 托管的外部表,表名称必须是格式
<catalog>.<schema>.<table>。 对于 Hive 元存储,它必须采用<schema>.<table>格式。运行完全刷新更新不会清理接收器中以前计算的结果数据。 这意味着任何重新处理的数据都会追加到接收器中,而不会修改现有数据。
不支持管道预期。
Serverless 出站控制仅支持 Kafka 和 Delta Lake Sink 连接器。 请参阅 什么是无服务器出口控制?。