数据工程的构建者模式:灵活基础

现代数据工程需要能够处理不断变化的数据源、转换逻辑和存储目的地的管道。 硬化的,单晶的管道设计往往导致在需求稍有变化时断裂的不规则系统。 构建器模式,一种既定的创建设计模式,提供了逐步构建复杂物体的结构化方法。 应用于数据管道,它可以使配置与执行脱钩,使工程师无需重写核心逻辑而调整管道。

理解构建器模式

起源和核心概念

构建器模式源于面向对象的编程,以解决构建对象与许多可选部分的问题。 与其使用一个具有众多参数或子分类的大型构建器来处理每个组合, 不如使用一个[ [FLT: 0]] 构建器[[[FLT: 1]] 对象提供逐个设置每个组件的方法。 最终的 [[FLT: 0] 方法组装了全部对象。 这种关切的分离使得构建过程可以在不同表达式中重新使用 。

分析:订购定制披萨

想想建筑师的模式,比如订购定制的比萨饼。 您指定了一块地壳、酱油、奶酪和顶层。 比萨师(厨师)知道如何将这些成分结合成一个完成的比萨饼。 同样的建筑师可以生产一个马格丽塔、夏威夷人或肉类爱好者的馅饼。 同样,一个数据管道建设师也可以从同一组建筑师的方法中组装不同的来源、转化和汇。

数据管道为何需要配置设计

数据管道很少是静态的。 从S3桶中摄入CSV文件并将其装入数据仓库的管道可能很快需要支持JSON、流源或额外的浓缩步骤。 没有可配置的设计,添加这些变化往往意味着复制和修改大量代码 — — 一种重复和错误的秘方。

  • 变换源系统: 从批量文件转移到事件流或切换数据库连接器.
  • 演化变换: 添加数据清洗,特征工程,或加入新的参考表.
  • 多功能目的地: 向多个数据存储库(如大查询,雪花,以及一个实时仪表板)写入结果,用于同一管道.
  • 测试和中转变体:[ 运行与开发与生产数据相同的逻辑,不改变代码.

构建器模式直接满足这些需求,让工程师 将管道进行申报[ ——定义哪些组件包括,如何连接,而基础组装逻辑保持不变.

可配置数据管道的核心组件

为了应用构建器模式,必须把数据管道分为离散的、可堆肥的构件。

数据来源

每个管道都以一个或多个来源开始:文件系统,数据库,流线平台(Kafka),API,或数据湖. 每个来源都有自己的配置(路径,证书,计划,投票间隔). 建构者可以提供诸如,,或等方法.

转换步骤

变换操作或丰富数据。 常见的例子包括过滤行、 解析嵌套 JSON、 汇总度量衡和加入数据集。 构建器方法如 、 和 允许工程师流畅地排序变换。

数据辛克斯

辛克斯是处理过的数据区域:关系数据库、云存储、消息队列或分析引擎。构建者可以支持多个汇,使用和,甚至允许链路将同样的数据发送到多个目的地。

连接器和中件

除了源和汇之外,管道通常需要处理错误、限制速度、计划验证器和监测钩。 这些交叉关注很容易作为构建者步骤添加,如或。

实施管道的构建器模式

典型的实现涉及一个管道构建器类,该类收集配置选项,以及一个构造方法[,该方法验证并返回一个完全构造的管道对象. 构造器暴露出将构建器本身还原用于链路的流畅方法.

class PipelineBuilder:
 def __init__(self):
 self._source = None
 self._transformations = []
 self._sinks = []
 self._retry_policy = None

 def with_source(self, source):
 self._source = source
 return self

 def add_transform(self, transform):
 self._transformations.append(transform)
 return self

 def add_sink(self, sink):
 self._sinks.append(sink)
 return self

 def with_retry(self, retry_policy):
 self._retry_policy = retry_policy
 return self

 def build(self):
 if not self._source or not self._sinks:
 raise ValueError("Source and at least one sink are required")
 return Pipeline(self._source, self._transformations, self._sinks, self._retry_policy)

利用建造者,管道的创建成为宣示:

pipeline = (PipelineBuilder()
 .with_source(S3CsvSource(bucket="data-landing", prefix="orders/"))
 .add_transform(FilterTransform(condition="status == 'active'"))
 .add_transform(AggregateTransform(group_by="customer_id", metrics=["sum(amount)"]))
 .add_sink(DatabaseSink(connection="prod_db", table="customer_orders"))
 .add_sink(ParquetSink(path="s3://analytics/orders/"))
 .with_retry(RetryPolicy(max_attempts=3, backoff_seconds=5))
 .build())

这种方法集中了配置,使得对装配和生产环境有不同参数的同一建构器易于再利用。

真实世界应用:建设弹性ETL管道

考虑一个电子商务公司,它需要从多个区域接收日常订购数据,清理和标准化数据,按类别计算每日收入,并将结果加载到一个报告数据库和一个数据湖中。它们利用构建者模式,创建了可重复使用的命令ETLBuilder

  1. Define source configs: 每个区域的命令来自不同的数据库(PostgreSQL, MySQL),但导出为共享的CSV格式. 构建器提供].
  2. 添加标准转换: 数据清理(删除无效的ID,验证货币代码)和浓缩(与产品目录合并以获得分类),这些通过和添加。
  3. 一组汇总:]。
  4. 冲向多个汇:[和]]。
  5. 建造并执行:[ 同一建造者可以先建造一条只读欧盟区域供测试的管道,然后将生产换到所有区域.

这种模式大大降低了代码重复:公司现在每个区域或环境维持一个构建者类,而不是多个特设脚本.

养恤金

  • 灵活性: 改变管道行为而不触碰执行逻辑。 需要添加新的转换吗? 只要用新的步骤调用 。
  • 保持性:管道定义读起来像高级食谱。 每个组件的配置都是孤立的,这使得调试和代码审查变得简单明了。
  • 续写: 构建器可以作为库进行包装. 团队在项目间重复使用同一构建器,只调整输入参数.
  • 伸缩性: 添加一个新的组件类型(如流槽)只需要延长建造器,而不是重写整个管道组装.
  • 试制性:[ 建构器可以用模拟源和汇来制造试验管,使单独的单元测试能够用于管道组装逻辑本身.

数据工程使用构建器模式的最佳做法

保留构建器纯配置

构建器只应收集并验证配置。 管道的实际执行应当由 [[FLT: 0]] 构造的 [[[FLT: 1] 对象负责。 这种分离使构建器保持简单和可测试性。

提前验证, 快速失败

在 ] 方法中, 验证所有所需的组件都存在, 并且配置一致( 如转换步骤参考已有源列 ) 。 丢弃描述性错误, 使用户确切知道缺少什么 。

利用不可移动的构造

在被称作之后,可以重置或再利用该建造器,以创建另一个具有不同环境的管道。避免在建构之间持续存在的存储状态,除非是故意的。

提供感知默认

对于复试策略或记录等可选组件,在建构器的构造器中设置合理的默认值。这可以将锅炉板最小化,同时允许覆盖。

版本您的构建器与管道

随着您数据基础设施的发展,构建者的API也会被开发。 Tag 构建者会在版本控制中释放,这样管道定义就可以连接到特定的构建者版本,防止在传播过程中意外地突破变化。

使用外部引用的复杂组件

对于内部细节很多的组件(例如Spark会话配置或自定义的UDF),考虑将它们作为预先建造的对象传递,而不是在管道构建器内建构它们. Refacting.Guru的构建者模式描述[为理解这种分离提供了极好的基础.

结论

构建者模式为数据工程团队创造了一个既强大又适应性的管道。 通过将什么(配置)与如何(执行)区分开来,它减少了技术债务,加速了对不断变化的商业需求的反应。 随着数据生态系统的复杂性 — — 实时流、多云存储和机器学习管道 — — 不断增长,构建者模式仍然是管理这一复杂性的可靠工具,而不会牺牲清晰度。

在设计下一条数据管道时,考虑采用构建器方法。它可能最初会觉得是一层额外的抽象,但灵活性和可维护性的长期收益远远大于前期成本。 对于数据工程的设计模式, Martin Fowler的分布式系统[提供了数据基础设施结构的更广泛视角。