什么是 Lakeflow Spark 声明性管道

Lakeflow Spark 声明性管道(SDP)是一个声明性框架,用于在 SQL 和Python中生成批处理和流式处理数据管道。 其核心概念包括管道、流、流式表、物化视图和接收端,它们协同工作,通过自动编排和增量更新来处理数据。

注释

Lakeflow Spark 的声明式管道需订购 高级版。 有关详细信息,请联系 Databricks 帐户团队。

什么是 SDP?

Lakeflow Spark 声明性管道是一个声明性框架,用于在 SQL 和 Python 中开发和运行批处理和流式处理数据管道。 Lakeflow SDP 扩展并可与 Apache Spark 声明性管道互作,同时在性能优化的 Databricks 运行时上运行,Lakeflow Spark 声明性管道 flows API 使用与 Apache Spark 和结构化流相同的数据帧 API。

SDP 的常见用例包括:

  • 从云存储(Amazon S3、Azure ADLS Gen2、Google Cloud Storage)和消息总线(Apache Kafka、Amazon Kinesis、Google Pub/Sub、Azure EventHub 和 Apache Pulsar)等源进行增量数据引入。
  • 使用无状态和有状态运算符的增量批处理和流式转换。
  • 事务存储(如消息总线和数据库)之间的实时流处理

有关声明性数据处理的更多详细信息,请参阅 Databricks 中的过程与声明性数据处理

SDP 有什么好处?

与使用 Apache SparkSpark Structured Streaming API 开发数据流程并通过 Lakeflow Jobs 手动编排在 Databricks Runtime 上运行相比,SDP 的声明性特性提供了以下优势。

  • 自动业务流程:SDP 自动协调处理步骤(称为“流”),以确保正确的执行顺序和最大并行度级别以实现最佳性能。 此外,管道会自动高效地重试暂时性故障。 重试过程从最精细且经济高效的单元开始:Spark 任务。 如果任务级重试失败,SDP 会继续重试流,然后在必要时重试整个管道。
  • 声明性处理:SDP 提供声明性函数,可将数百甚至数千行手动 Spark 和结构化流式处理代码减少为几行。 SDP AUTO CDC API 简化了变更数据捕获(CDC)事件的处理,同时支持 SCD 类型 1 和 SCD 类型 2。 它无需手动代码来处理无序事件,并且不需要了解流式处理语义或水印等概念。
  • 增量处理:SDP 为具体化视图提供增量处理引擎。 若要使用它,请使用批处理语义编写转换逻辑,并且引擎只会尽可能处理数据源中的新数据和更改。 当源中发生新数据或更改时,增量处理会降低低效的重新处理,并且无需手动代码来处理增量处理。

重要概念

下图演示了 Lakeflow Spark 声明性管道最重要的概念。

此图显示了 SDP 的核心概念在非常高的水平上如何相互关联

数据集

管道生成三种类型的数据集,每个数据集具有不同的处理语义:

数据集类型 如何处理记录
流式处理表 假设只追加源,每个记录都处理一次。 流式处理表适用于持续增长数据的引入和增量处理。
具体化视图 根据需要重新计算结果以反映数据的当前状态。 物化视图适合用于转换、聚合,或预先计算供多个下游数据集使用的结果。
查看 按需评估,不持久化。 对于无需发布到目录中的中间转换和检查操作,请使用视图。

流式处理表是 Unity 目录托管表的一种形式,也是流式处理目标。 流式处理表可以写入其中一个或多个流式处理流(追加AUTO CDC)。 您可以显式定义流式流,并将其与目标流式表分开定义;也可以将其作为流式表定义的一部分进行隐式定义。

具体化视图也是 Unity 目录托管表的一种形式,是批处理目标。 具体化视图可以有一个或多个具体化视图流被写入其中。 具体化视图不同于流式处理表,因为您总是将流隐式定义为具体化视图定义的一部分。

有关详细信息,请参阅 流式处理表具体化视图

何时使用视图、具体化视图和流式表

实现管道查询时,请选择最适合用例的数据集类型。

考虑使用视图来:

  • 将大型查询或复杂查询分解为易于管理的查询。
  • 使用期望来验证中间结果。
  • 减少不需要保留的结果的存储和计算成本。 由于表已具体化,因此它们需要额外的计算和存储资源。

在以下情况下考虑使用具体化视图:

  • 多个下游查询使用表。 由于视图是按需计算的,因此每次查询视图时都会重新计算视图。
  • 其他管道、作业或查询将使用表。 由于视图未具体化,因此只能在同一管道中使用它们。
  • 你想要在开发过程中检查查询的结果。 由于表是具体化的,可以在管道外部查询,因此在开发过程中使用表有助于验证计算的正确性。 验证后,将不需要具体化的查询转换为视图。

在以下情况下考虑使用流式表:

  • 查询是针对持续或以增量方式增长的数据源定义的。
  • 应以增量方式计算查询结果。
  • 管道需要高吞吐量和低延迟。

注释

流数据表始终是基于流数据源定义的。 你还可以将流式处理源与 AUTO CDC ... INTO 结合使用以应用 CDC 源中的更新。 请参阅 AUTO CDC API:使用管道简化变更数据捕获

Flows

流是 SDP 中支持流式处理和批处理语义的基础数据处理概念。 流从源读取数据,应用用户定义的处理逻辑,并将结果写入目标。 SDP 与 Spark 结构化流共享相同的流式处理类型(追加更新完成)。 (目前,仅公开 追加更新 流。有关详细信息,请参阅 结构化流式处理中的输出模式

Lakeflow Spark 声明性管道还提供其他流类型:

  • AUTO CDC 是 Lakeflow SDP 中唯一的流式处理流,可处理无序 CDC 事件,并支持 SCD 类型 1 和 SCD 类型 2。 自动 CDC 在 Apache Spark 声明性管道中不可用。
  • 具体化视图 是 SDP 中的批处理流,仅尽可能处理源表中的新数据和更改。

有关详细信息,请参阅 使用 Lakeflow Spark 声明性管道流以增量方式加载和处理数据

Sinks

接收端是管道的流处理目标,支持 Delta 表、Apache Kafka 主题、Azure EventHubs 主题和自定义 Python 数据源。 接收器可以写入一个或多个流式流(追加更新)。

有关更多详细信息,请参阅 Lakeflow Spark 声明式管道中的输出端

Pipelines

管道是 Lakeflow Spark 声明性管道中的开发和执行的单元,是定义的流、流式处理表、具体化视图和接收器的容器。 通过在管道源代码中定义这些对象,然后运行管道,来使用 SDP。 管道运行时,它会分析定义的对象的依赖项,并自动协调其执行顺序和并行化。

有关详细信息,请参阅什么是管道?

数据引入

管道支持 Azure Databricks 中提供的所有数据源。 Databricks 建议为大多数引入用例使用流式处理表。 对于云对象存储中的文件,自动加载程序提供增量的幂等加载。 对于流数据,管道可以直接从消息总线(如 Apache Kafka、Azure 事件中心、Amazon Kinesis 和 Google Pub/Sub)引入。 请参阅 在管道中加载数据

数据质量

期望规则是数据集上的可选规则,用于在数据流经管道时验证数据。 将预期定义为 SQL 布尔约束,并指定记录失败时会发生什么情况:警告、删除记录或更新失败。 请参阅通过管道预期管理数据质量

Delta 集成

由管道创建和管理的所有表均为 Delta 表。 它们具备与 Delta Lake 相同的保障,包括 ACID 事务、时间回溯和架构强制。 请参阅 Azure Databricks 中的 Delta Lake 是什么?

其他资源