本页解释了如何在数据管道生命周期中使用Lakeflow管道,从最初的设计决策到大规模运行,以及每个阶段背后的权衡。 每个章节都链接到教你如何操作的文章。
本指南假设对核心数据工程概念有熟悉度。 如果你是管道新手,可以先从 Apache Spark声明式管道 开始,了解产品内容及其背后的声明式模型,然后再学习 教程:使用变更数据捕获构建ETL管道。
管道生命周期概述
管道通过六个阶段:
- 规划与设计:决定你要建造什么,选择合适的工具、语言和计算。
- 数据导入:可靠且逐步地将源数据导入流水线。
- 转换与建模:清理、验证、连接并塑造数据,形成消费者可信的表格。
- 运营化:将流水线置于版本控制之下,测试、调度并在各环境中推广。
- 在生产环境中运行:在流水线无人值守运行时,监控、提醒、调试、回填、保护并跟踪沿革。
- 走向成熟并实现规模化:确认已做好生产就绪准备,并在业务量和团队规模增长的同时保持流水线的健康运行。
这些阶段并非严格顺序,但它们对应出题的顺序。 由于Lakeflow流水线处理编排、检查点、重试和增量处理,你在每个阶段的工作更多是设计决策,而非实现。
规划和设计
你的第一个决定会影响后续的一切。 关于声明式模型与自己编写程序步骤的区别,请参见 Azure Databricks 中的程序式与声明式数据处理。
有几个选择决定了你的起始配置:
- 一个独立数据集或一个管道。 单个具体化视图或流表可以在SQL中定义为独立数据集,Azure Databricks负责管理其背后的刷新流水线。 当您需要使用 Python 编写、配置输出端或进行多阶段编排时,可以将 Lakeflow 流水线作为一个整体来编写和运行。 参见 独立管道与湖流量管道。
- SQL或Python(或者两者兼用)。 SQL 适合主要由过滤、连接和聚合组成的转换。 Python 适合自定义逻辑、外部库,或通过程序生成许多类似的表。 这种选择是以文件为单位进行的,而不是针对整个流水线,因此你可以将两种方式混用,而不必一开始就决定。
- 无服务器或经典计算。 推荐默认是无服务器配置,并且会移除集群配置。 当你需要特定实例类型、自定义集群策略或初始化脚本时,选择经典版。 参见 “配置无服务器流水线 ”和 “配置流水线经典计算”。
- 触发式执行或连续执行。 启动触发,因为它运行时只消耗计算量。 连续模式能以最小延迟保持计算运行以处理新数据,这通常是成本最高的因素,因此应保留给经过验证的延迟需求。 请参阅触发与连续管道模式。
流水线是根据代码引用的数据集推断执行图,所以设计工作主要是给数据集命名和排序。 核心决策是每个输出应采用哪种类型:用于大量附加增量数据的 流表 ,还是用于重新计算的聚合和连接的 具体视图 。 这种选择决定了成本和正确性,因为增量处理会随着新数据的增长速度增长,而完整的重新计算则会随着你的整个历史记录而扩展。 关于哪种类型适合哪种岗位,请参见 “什么是管道?”。
因为流水线代码是普通的Python和SQL,你可以在自己的编辑器里编写、lint并验证它,然后再部署到共享工作区。
在本阶段
在这个阶段需要思考的问题:
- 我该如何在独立数据集和完整流水线之间做出选择?
- 我该如何识别我的数据源并弄清楚如何连接它们?
- 在编写任何代码之前,我该如何设计我的流水线架构?
- 我该如何选择文件格式和存储层?
- 我该如何搭建本地开发环境?
- 在开始建造之前,我该如何规划规模和估算成本?
引入数据
设计中的核心问题在于,数据源究竟是只追加,还是就地更改。 这决定了你如何建模目标:
- 仅附加的源,如云存储中的文件或消息总线上的事件,会被导入流 式表,流式表会检查其进度,因此重启既不会重新处理也不会丢弃数据。 Auto Loader 处理文件,在新文件到达时发现它们,并推断和演进架构。 Apache Kafka、Azure 事件中心、Amazon Kinesis 和 Google Pub/Sub 等消息总线中的数据可直接读入流式表。 在下游进行去重,因为总线可能会多次投递同一事件。 对于 Azure 事件中心,请参见 使用 Azure 事件中心 作为管道数据源。
-
负责更新和删除行的来源,如大多数数据库和许多软件即服务(SaaS)系统,使用变更数据捕获(CDC)。 每次运行都进行完整复制不仅浪费资源,而且源数据越大,复制速度就越慢,因此,CDC 只读取自上次运行以来发生变化的行。
AUTO CDCAPI 无需手写合并逻辑即可应用这些更改;参见 AUTO CDC API:借助管道简化变更数据捕获。 一个 流会将 CDC 应用于流式表,多个流可以为同一个表提供数据,这就是将多个源扇入到单个目标中的方式。
检查点和重试是自动的,因此流水线会从最后处理的偏移量开始,而不是全部重新处理。 有两种保障措施是选择加入的:
- 救援数据列捕捉与预期模式不匹配的记录。
- 期望 应用你定义的行级动作。
如果流式检查点失效,优先选择成本最低且能保留表数据的恢复方式。
在本阶段
在这个阶段需要思考的问题:
转换与建模
转换将导入的数据转化为人们和工具可信赖的干净表格。 这就是奖 章 图案(青铜到银再到金)具体成形的地方。
清洁和确认是第一位的。 预期是 Lakeflow 管道的一项内置功能:它是管道对每次运行中每一行数据都会评估的数据质量约束,并报告通过和失败的计数,因此数据质量是持续保障的,而不是一次性检查。 决定当某一行处理失败时应如何处理(发出警告并保留该行、丢弃该行,或使更新失败),以及关口应设置在哪里。 关卡通常位于青铜级与白银级之间的分界处,因此下游环节的所有内容都可视为可信,无需再次核查。
连接和聚合形成了银到金的阶梯。 实体化视图适合对现有表进行批处理式的连接或聚合,因为它保持结果与源头一致:当查询和源允许时,它会增量刷新,否则则完整重新计算,结果无论如何都相同。 这使得当正确性比延迟更重要时,它是正确的选择,因为它在维度变化时重新计算连接。 详见 “如何刷新管道?”。 对实时流进行连接会产生无界状态,因此流式连接和聚合需要借助水印来限定流水线等待迟到数据的时长。
有两个正确性原则贯穿这一阶段:
-
幂零性 意味着一个管道无论对同一输入运行多少次,都会产生相同的结果。 Lakeflow 管道对其管理的部分(如检查点读取和基于
AUTO CDC键的上源)具有幂等性;通过避免在重新计算视图中使用非确定性函数,保持自身逻辑幂等性。 - 至少一次处理与恰好一次处理。 管理的Delta到Delta表会将每个微批次的输入和输出一起提交,默认只提交一次。 不过,这种保证在一些边界场景下就不再适用,例如自定义接收端、非 Delta 目标,或未经验证的自定义源;在这些情况下,你需要将写入视为至少一次,并使写入具备幂等性,例如基于某个键执行 upsert。
缓慢变化维度(SCD)也存在于此: AUTO CDC 直接实现了SCD类型1和类型2,因此您只需设置类型,而无需编写历史跟踪逻辑。
在本阶段
在这个阶段需要思考的问题:
- 我该如何清理和验证输入数据?
- 我如何在维度缓慢变化(SCD)下追踪历史?什么是SCD?
- 我该如何连接流媒体和静态数据?我如何高效地聚合数据?
- 我如何建模数据以便下游使用?
- 我如何确保Lakeflow管道的加工保证?
- 至少一次处理与恰好一次处理:有什么区别,我需要哪一种?
- 我该如何处理迟到或乱序的数据?
可操作化
运营化是将流水线从仅供你个人运行的流程,转变为团队能够可重复地构建、测试和发布的流程。 流水线是源代码加上配置,因此适用普通的软件工程实践。
要按预定计划运行流水线,请将其置于在工作流中运行流水线:Databricks 建议使用作业来调度和编排流水线,这样你还可以将流水线与其他任务协调起来,例如串联下游报表或多个流水线。 在一次运行中, 流水线 会下单并行处理自己的数据集,因此编排只协调流水线外的任务。
在本阶段
在这个阶段需要思考的问题:
在生产环境中运行
一旦流水线开始在真实数据上无人值守运行,工作重点就变成了判断其运行是否正常,并在出现问题时进行修复。
监测分为三个深度层级。 作业和流水线列表可一目了然地显示最近运行的状态。 流水线监控界面会按状态颜色区分每个表格和流程,包括行数、数据质量指标和流表的待办项指标。 两者下方的事件日志是任何程序化或历史性内容的真实来源。 配置故障通知,让你在利益相关者报告之前就知道运行异常。 有关监测表面的概述,请参见 “监测管道”。
从图中高亮的失败点倒推到事件日志中的完整错误细节进行调试,然后只重跑失败的部分。 重试行为因触发器而异:手动触发的更新会禁用自动重试,这样你会立即看到错误,而计划更新会重试可恢复的失败。 因此,生产警报在重试后可能会自行解除,而同样的故障在交互式开发时则不会。
将回填建模为一个显式的一次性流程,使其流向与常规增量流程相同的目标。 将其单独分离出来,可以记录历史是何时以及如何加载的,同时使稳态逻辑保持简单。
通过控制谁可以操作流水线、作为专用服务主体而非个人账户运行,以及将凭证保存在秘密作用域而非源代码中来保障流水线安全。 血统是自动的,记录到列级。 管道通过 Lakeflow 管道中的接收器 将数据写入外部系统,而这正是上述“至少一次”语义适用的边界环节。
在本阶段
在这个阶段需要思考的问题:
成熟度与规模
成熟的流水线可在无人值守的情况下运行,并且无需重写也能持续扩展。 确认是否已准备就绪并规划如何实现规模化,是这一阶段的核心。
生产准备度是一份涵盖数据质量、可靠性、可观察性、部署、成本和治理的清单。 把每一个未勾选的项都当作已知的漏洞:每个可能接收坏数据的数据集是否有预期,流水线是调度的还是手动启动的,失败通知是否已配置,是否作为服务主体运行,是否从至少开发和生产目标的版本控制部署。 数据质量检查和告警是最容易添加且成本最低的,也最有可能发现未被察觉的异常运行。
在流水线健康状况出现恶化的明确信号时进行扩容:
- 更新时间正在上升。
- 自动缩放屡次触及上限。
- 成本增长速度快于基础业务。
- 实体化视图正在回退到完全重新计算。
先尝试计算层面的杠杆,比如转向无服务器,或者根据你的延迟需求匹配其性能模式。 除此之外,如何组织跨管道的数据集最为重要:
- 流水线有并发限制:它同时只更新一定数量的数据集。 一旦管道的数据集超过该限制,额外的更新会在队列中等待,因此管道的总更新时间会增加。
- 将相关数据集分组,并拆分无关数据集。 按域、相同的刷新频率和依赖关系分组;按所有权、层级和延迟边界拆分。 例如,将数据摄入与转换分离,可以避免缓慢的摄入拖慢下游所有环节,同时使每个管道都足够小,从而保持在并发限制之内。
合并两条小型管道比拆分一个已经在生产的大型管道要容易得多。 有关如何对数据集进行分组和拆分,请参见 跨 Lakeflow 管道组织数据集。
在本阶段
在这个阶段需要思考的问题: