在管道中使用汇聚器

使用 Lakeflow 管道 sink API 和 ,将经管道转换后的记录写入外部数据接收端。 外部数据接收器包括 Unity 目录托管表和外部表,以及 Apache Kafka 或 Azure 事件中心等事件流服务。 还可以通过编写 Python 代码,将数据接收器用于写入自定义数据源。

有关接收器概念的概述以及何时使用它们,请参阅 Lakeflow 管道中的接收器

注释

接收端工作流

当事件数据从流式数据源引入管道时,可以通过管道转换处理和优化该数据。 然后使用追加流处理过程将转换后的数据记录传输到接收器。 使用 create_sink() 函数创建此接收器。 有关 create_sink 函数的更多详细信息,请参阅 Sink API 参考

如果您的管道创建或处理流式事件数据,并准备好数据记录供写入到接收端,那么您就可以使用数据接收端。

实现接收器包括两个主要步骤:

  1. 创建数据汇聚器。
  2. 使用 追加流更新流 将准备好的记录写入接收器。

创建下沉器

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 参数是必需的。 其他参数(如 keypartitionheaderstopic)是可选的。

有关 append_flow 装饰器的更多信息,请参阅 默认流和追加流

局限性

  • 仅支持 Python API。 不支持 SQL。

  • 仅支持流式处理查询。 不支持批处理查询。

  • 只有 append_flowupdate_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 连接器。 请参阅 什么是无服务器出口控制?

其他资源