什么是 Lakeflow 管道?

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

Lakeflow 管道是对 Apache Spark™ 声明性管道(SDP)的扩展。 若要详细了解 SDP 及其与 Lakeflow 管道的比较方式,请参阅 Apache Spark 声明性管道

注释

Lakeflow 管道需要使用高级套餐。 有关详细信息,请联系 Databricks 帐户团队。

管道的优点是什么?

与在 Databricks Runtime 上通过 Lakeflow Jobs 进行手动编排,使用 Apache SparkSpark Structured Streaming API 开发数据工程流程相比,管道的声明式特性具有以下优势:

  • 自动业务流程:管道按正确的顺序运行处理步骤(称为“流”),并按最大并行度逐步重试暂时性故障,从 Spark 任务到流到整个管道。
  • 声明性处理:声明性函数将数百行手动 Spark 和结构化流式处理代码减少到几个行。 AUTO CDC API 处理变更数据捕获(CDC)事件(包括 SCD 类型 1 和类型 2),无需手动代码处理无序事件或流式处理概念(如水印)。
  • 增量处理增量处理 引擎使具体化视图保持最新状态:使用批处理语义编写转换逻辑,并且引擎尽可能仅重新处理新的或更改的源数据。

重要概念

下图说明了管道最重要的概念。

此图显示了管道的核心概念如何在非常高的级别相互关联

数据集

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

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

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

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

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

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

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

考虑使用视图来:

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

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

  • 多个下游查询使用表。 由于物化视图会缓存查询结果,因此后续查询会直接读取预先计算的结果,而不是在每次访问时重新执行该查询。
  • 其他管道、作业或查询将使用表。 由于物化视图被物化为 Unity Catalog 表,因此定义该视图的管道之外的用户也可以查询该视图。 视图不会被物化,因此只能在同一管道内使用。
  • 你想要在开发过程中检查查询的结果。 由于物化视图会被实际存储,并且可以在管道之外进行查询,因此你可以在开发过程中验证计算结果的正确性。 验证后,将不需要具体化的查询转换为视图。
  • 查询执行聚合或联接,或者源数据可能会因更新和删除而更改,而不只是增长。 物化视图会使其结果与源数据的当前状态保持一致,而流式表则专为仅追加型数据源设计,并且每条记录只处理一次。

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

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

注释

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

Flows

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

管道还提供其他流类型:

  • AUTO CDC 是 Lakeflow 管道中唯一的流式处理流,用于处理无序 CDC 事件并支持 SCD 类型 1 和 SCD 类型 2。 自动 CDC 在 SDP 中不可用。
  • 物化视图是管道中的一种批处理流,在可能的情况下,它只处理源表中的新数据和变更。

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

Sinks

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

有关详细信息,请参阅 Lakeflow 管道中的接收器

Pipelines

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

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

你还可以在 Lakeflow 管道之外定义独立的物化视图和流式表,在这种情况下,Azure Databricks 会为你管理该管道。 若要比较这两种方法,请参阅 独立管道与 Lakeflow 管道

数据引入

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

数据质量

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

Delta 集成

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

其他资源