0%

数据管道全链路设计范式

数据管道全链路设计范式

一句话定位:数据管道不是”一串脚本”,而是一套”分层解耦 + 算子可复用 + 可重跑 + 资源隔离”的工程系统。面试要把”接入→清洗→转换→输出”讲成可治理的范式。

1. 分层解耦:接入 → 清洗 → 转换 → 输出

  • 每层职责单一,通过标准数据契约(schema / 格式)衔接;
  • 单层可独立替换、独立扩缩(如换 source 不影响 transform);
  • 故障定位时按层切分,避免”一处坏全线瘫”。

2. 算子抽象与可复用性

  • 把处理单元抽象为统一接口的算子(map / filter / join / aggregate;Flink 算子、Spark 算子、Ray Actor);
  • 沉淀为可复用算子库 + 编排模板,业务逻辑靠组合而非重写;
  • 对应”平台化意识”:一次优化沉淀为通用能力。

3. 幂等与重试设计

  • 幂等:同一输入多次执行结果一致(唯一键 / 去重表 / upsert / 幂等写入),使重试安全;
  • 重试 + 断点续跑:失败任务可重跑(至少一次 → 配合幂等达成”有效一次”),靠位点/offset/checkpoint 续跑,避免全量重算;
  • 对应 1.1.4 Flink 的 Exactly-Once 思想。

4. CPU 密集与 GPU 密集的资源池分离

  • 解析/清洗/特征(CPU 密集)与推理/打标(GPU 密集)分属不同资源池(如 CPU 集群 vs GPU 集群);
  • 用 Ray / 调度器做异构编排(如 head 在 CPU 集群、worker 在 GPU 集群),避免互相抢占;
  • 对应 1.1.5 Ray 的异构资源调度优势。

5. 可观测与背压

  • 监控各环节吞吐/延迟/积压;
  • 背压(Backpressure):下游过载时反向限速,防止雪崩(Flink 原生背压机制)。

参考:Google《Data Pipelines》相关实践;Flink / Spark Structured Streaming 背压机制文档