M8-019M8: ML Systems, Engineering & ResearchData Pipelines & StreamingEasy
Mastery:

Data Pipelines & Streaming: 解释批处理与流处理的差异与选型。

📐 Mathematical Definition
batch: window=hours/days;stream: window=seconds\text{batch}:\ \text{window}=\text{hours/days};\qquad \text{stream}:\ \text{window}=\text{seconds}
⚡ Executive Summary
Core Concept: 批处理(高吞吐、高延迟、可重跑)vs 流处理(低延迟、逐条、状态管理复杂);按'延迟需求'选。

📌 Key Takeaways

  • •
    批处理:定时调度、全量/增量、高吞吐、可重跑
  • •
    流处理:持续消费、逐条/微批、低延迟、状态管理
  • •
    选型:延迟需求(秒级 vs 小时级)+ 复杂度 + 成本

📐 Mathematical Derivations

数学机理:<strong>批处理与流处理</strong>——(1) <strong>批处理(batch)</strong>——(a) <strong>模式</strong>——定时调度(小时/天),处理<strong>有界数据</strong>(一个时间窗的数据);(b) <strong>优点</strong>——(i) <strong>高吞吐</strong>(可用全部资源处理大批数据);(ii) <strong>可重跑</strong>(数据是有界的,可重算);(iii) <strong>简单</strong>(无状态管理);(iv) <strong>成本低</strong>(按需启动);(c) <strong>缺点</strong>——(i) <strong>延迟高</strong>(小时/天级);(ii) 无法处理'实时'需求。(2) <strong>流处理(stream)</strong>——(a) <strong>模式</strong>——持续消费<strong>无界数据流</strong>,逐条(或微批)处理;(b) <strong>优点</strong>——(i) <strong>低延迟</strong>(秒/亚秒级);(ii) 适合'实时'需求(风控/推荐/监控);(c) <strong>缺点</strong>——(i) <strong>状态管理复杂</strong>(窗口状态、去重状态);(ii) <strong>乱序与迟到数据</strong>(需 watermark 机制);(iii) <strong>精确一次语义难</strong>(需 checkpoint + 两阶段提交);(iv) <strong>运维复杂</strong>(常驻服务);(v) <strong>成本高</strong>(常驻资源)。(3) <strong>选型依据</strong>——(a) <strong>延迟需求</strong>——(i) 秒级 → 流处理;(ii) 小时/天级 → 批处理;(b) <strong>数据量</strong>——(i) 海量 → 批处理(吞吐优势);(ii) 中等 → 流处理;(c) <strong>复杂度容忍</strong>——(i) 简单优先 → 批处理;(d) <strong>成本</strong>——(i) 常驻成本 vs 按需成本;(e) <strong>一致性要求</strong>——(i) 精确一次 → 流处理需专门设计。(4) <strong>关键技术(流处理)</strong>——(a) <strong>窗口(window)</strong>——(i) <strong>滚动窗口</strong>(不重叠);(ii) <strong>滑动窗口</strong>(重叠);(iii) <strong>会话窗口</strong>(按活动间隔);(b) <strong>Watermark</strong>——处理'乱序与迟到数据'('水印'表示'该时间点之前的数据应该都到了');(c) <strong>状态存储</strong>(RocksDB/内存——支持大状态);(d) <strong>Checkpoint</strong>(容错——定期快照状态);(e) <strong>精确一次</strong>(端到端——需 sink 支持事务);(f) <strong>背压(backpressure)</strong>(消费慢于生产时的处理)。(5) <strong>流批一体</strong>——(a) <strong>Lambda 架构</strong>——批处理层(准确)+ 速度层(低延迟)+ 服务层(合并);<strong>缺点</strong>——两套代码(逻辑重复);(b) <strong>Kappa 架构</strong>——只用流处理(重放历史数据代替批处理);<strong>优点</strong>——一套代码;<strong>缺点</strong>——历史重放成本高;(c) <strong>流批一体引擎</strong>(Flink/Spark Structured Streaming)——同一套 API 处理批与流。(6) <strong>中间形态</strong>——(a) <strong>微批(micro-batch)</strong>(Spark Streaming——秒级延迟,简单);(b) <strong>近线(nearline)</strong>(分钟级——用批处理框架跑高频任务)。<strong>实践建议</strong>——(a) <strong>延迟需求决定</strong>(秒级→流、小时级→批);(b) <strong>流处理注意状态/乱序/精确一次</strong>;(c) <strong>优先'流批一体'</strong>(避免两套代码);(d) <strong>不要为了'实时'而过度设计</strong>(很多场景分钟级足够)。<strong>度量</strong>——(a) 端到端延迟;(b) 吞吐;(c) 状态大小;(d) 成本。

🏭 Production Trade-offs

深度剖析与工程权衡:① <strong>'延迟需求决定选型'</strong>——很多场景'分钟级'足够,不必上流处理;面试中能指出'不要过度设计'是深度理解的标志。② <strong>'流处理的状态管理'是主要复杂度</strong>——窗口状态/去重状态/checkpoint。③ <strong>'乱序与迟到数据'需 watermark</strong>——这是流处理的核心难点。④ <strong>'Lambda 的两套代码'是痛点</strong>——故有 Kappa 与流批一体。⑤ <strong>'微批是实用折中'</strong>——秒级延迟 + 简单实现(Spark Streaming)。⑥ <strong>面试要点</strong>——被问'批处理还是流处理',应给出'<strong>对比(延迟/吞吐/复杂度/成本)+ 选型依据(延迟需求)+ 流处理关键技术(窗口/watermark/状态/checkpoint/精确一次)+ 流批一体</strong>';能指出'不要过度设计'是深度理解的标志。
⚠️ Common Interview Pitfalls
  • ✕
    为'分钟级'需求上流处理(过度设计)
  • ✕
    流处理忽略乱序与迟到数据
🎯 Interviewer Follow-ups
  • ?
    什么场景必须用流处理?
  • ?
    流处理的'状态管理'难在哪?
📚

Associated Knowledge Base Guides & Mindmaps

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

← PreviousM8-018: Data Pipelines & Streaming: 描述一个典型离线数据管道。📋Back to BankNext →M8-020: Data Pipelines & Streaming: 解释数据质量监控的维度与手段。