M8-024M8: ML Systems, Engineering & ResearchData Pipelines & StreamingHard
Mastery:

Data Pipelines & Streaming: 解释流批一体的架构(Lambda vs Kappa)。

📐 Mathematical Definition
Lambda: batch+speed+serving;Kappa: stream only+replay\text{Lambda}:\ \text{batch}+\text{speed}+\text{serving};\qquad \text{Kappa}:\ \text{stream only}+\text{replay}
⚡ Executive Summary
Core Concept: Lambda:批处理层(准确)+ 速度层(低延迟)+ 服务层;Kappa:只用流处理(重放历史);流批一体引擎统一两者。

📌 Key Takeaways

  • •
    Lambda:两套代码(批+流),逻辑重复是痛点
  • •
    Kappa:一套流代码,用'重放历史'替代批处理
  • •
    流批一体引擎(Flink/Spark SS):同一套 API 处理批与流

📐 Mathematical Derivations

数学机理:<strong>三种架构</strong>——(1) <strong>Lambda 架构</strong>——(a) <strong>三层</strong>——(i) <strong>批处理层(batch layer)</strong>——处理全量历史,产出'准确但延迟高'的视图;(ii) <strong>速度层(speed layer)</strong>——处理实时数据,产出'低延迟但可能近似'的视图;(iii) <strong>服务层(serving layer)</strong>——合并两者的查询结果;(b) <strong>优点</strong>——(i) 批处理保证准确性(可重算全量);(ii) 速度层保证低延迟;(c) <strong>痛点</strong>——(i) <strong>两套代码</strong>(批与流的逻辑需分别实现——<strong>最容易不一致</strong>);(ii) 维护成本高;(iii) 两套结果需对齐。(2) <strong>Kappa 架构</strong>——(a) <strong>核心思想</strong>——<strong>只用流处理</strong>;需要'重算历史'时,<strong>重放(replay)</strong>消息队列中的历史数据(而非用批处理);(b) <strong>优点</strong>——(i) <strong>一套代码</strong>(逻辑不重复);(ii) 一致性高;(c) <strong>痛点</strong>——(i) <strong>重放成本</strong>(重放数月/数年的数据成本高、耗时长);(ii) 消息队列需保留长历史(存储成本);(iii) 不适合'超大规模历史'的场景。(3) <strong>流批一体(unified)</strong>——(a) <strong>做法</strong>——用同一套 API/引擎处理批与流(Flink、Spark Structured Streaming、Beam);(b) <strong>核心抽象</strong>——'<strong>有界流</strong>'(批)与'<strong>无界流</strong>'(流)的统一;(c) <strong>优点</strong>——(i) 一套代码;(ii) 批处理作为'流的一个特例'(重放历史);(iii) 一致的语义(事件时间/watermark/状态);(d) <strong>现状</strong>——Flink 的'流批一体'是当前主流方向。<strong>关键概念(流处理)</strong>——(a) <strong>事件时间 vs 处理时间</strong>——事件时间是'数据产生的时间'(正确但需处理乱序);处理时间是'处理的时间'(简单但受延迟影响);(b) <strong>Watermark</strong>——'事件时间的水位'(表示'该时间之前的数据应该都到了');(c) <strong>窗口</strong>(滚动/滑动/会话);(d) <strong>状态</strong>(窗口状态/去重状态——RocksDB);(e) <strong>精确一次</strong>(checkpoint + 两阶段提交);(f) <strong>背压</strong>。<strong>选择依据</strong>——(a) <strong>团队小/简单</strong> → Kappa 或流批一体;(b) <strong>需要'全量重算'且历史极大</strong> → Lambda(批处理层处理历史);(c) <strong>新建系统</strong> → 流批一体引擎(避免 Lambda 的两套代码)。<strong>与其他问题的关系</strong>——(a) 与'批处理 vs 流处理'(上一题);(b) 与'训练-服务一致性'(两套代码易不一致);(c) 与'数据版本'(重放需要历史)。<strong>实践建议</strong>——(a) <strong>新建优先流批一体</strong>(避免两套代码);(b) <strong>Lambda 用于'历史极大 + 需低延迟'</strong>;(c) <strong>Kappa 需评估重放成本</strong>;(d) <strong>注意事件时间与 watermark</strong>;(e) <strong>精确一次需专门设计</strong>;(f) <strong>监控两套结果的一致性</strong>(Lambda)。<strong>度量</strong>——(a) 代码重复度;(b) 批/流结果的一致性;(c) 重放成本;(d) 端到端延迟。

🏭 Production Trade-offs

深度剖析与工程权衡:① <strong>'Lambda 的两套代码易不一致'是核心痛点</strong>——面试中能指出是深度理解的标志。② <strong>'Kappa 的重放成本'是主要限制</strong>——历史极大时不可行。③ <strong>'流批一体是当前方向'</strong>——Flink/Spark SS 统一 API。④ <strong>'事件时间 vs 处理时间'是流处理的关键区分</strong>——影响正确性。⑤ <strong>'Watermark'处理乱序</strong>——流处理的核心机制。⑥ <strong>面试要点</strong>——被问'Lambda vs Kappa',应给出'<strong>Lambda(批+速度+服务,两套代码)/ Kappa(只流+重放)/ 流批一体(统一 API)+ 关键概念(事件时间/watermark/状态/精确一次)</strong>';能指出'两套代码易不一致'是深度理解的标志。
⚠️ Common Interview Pitfalls
  • ✕
    Lambda 的批与流逻辑不一致(结果对不上)
  • ✕
    Kappa 不评估重放成本(历史极大时不可行)
🎯 Interviewer Follow-ups
  • ?
    Lambda 的'两套代码'为什么危险?
  • ?
    Kappa 的'重放'成本如何控制?
📚

Associated Knowledge Base Guides & Mindmaps

Explore the comprehensive technical article, exam cards, and global architecture tree.

← PreviousM8-023: Data Pipelines & Streaming: 解释数据版本管理与快照。📋Back to BankNext →M8-025: Training Platforms & Experiment Tracking: 解释实验管理需要记录哪些信息。