惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Apple Machine Learning Research
Apple Machine Learning Research
博客园_首页
G
Google Developers Blog
aimingoo的专栏
aimingoo的专栏
罗磊的独立博客
博客园 - 【当耐特】
M
MIT News - Artificial intelligence
D
Docker
博客园 - 三生石上(FineUI控件)
博客园 - 司徒正美
人人都是产品经理
人人都是产品经理
博客园 - 叶小钗
月光博客
月光博客
S
SegmentFault 最新的问题
Jina AI
Jina AI
Blog — PlanetScale
Blog — PlanetScale
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
博客园 - Franky
L
LangChain Blog
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
Microsoft Azure Blog
Microsoft Azure Blog
阮一峰的网络日志
阮一峰的网络日志
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
Last Week in AI
Last Week in AI

暗无天日

读:AI Agent 安全日志——从可见性与隐私的两难说起 - 暗无天日 AI写作的语言指纹——如何让文字不那么像机器 - 暗无天日 读:50 条 Claude Code 技巧——一个工程经理的六个月使用心得 读:AI 辅助开发为什么让 E2E 测试更有价值 - 暗无天日 读:在Emacs中使用Claude Code(Spacemacs适配版) - 暗无天日 Claude Code 背后的工程哲学——读 Agent Harness Engineering 读:Agent Harness Engineering——AI 智能体不只是模型,还有套件 - 暗无天日 browser-harness:让 AI 直接接管你的浏览器 - 暗无天日 读:Security-First CI/CD —— DevSecOps 自动化实践指南 TIL: 数字小键盘的小数点陷阱与行内算术求值 - 暗无天日 读:Immutability 不是万能药,它是一种权衡 - 暗无天日 Conducty:给 Claude Code 加上项目记忆和并行执行能力 - 暗无天日 读 — GitHub Trending 里的 Claude Code 技能包 读 — Prompt Caching 省钱指南 TIL: Emacs 中那些跟鼠标配合的冷门快捷键 - 暗无天日 读:Anvil——把 Emacs 变成 AI 的工具服务器 读:Emacs 代码折叠终极指南 - 暗无天日 读:Clojure 搭车客指南 - 暗无天日 git推送失败后恢复仓库损坏的完整记录 - 暗无天日 多智能体系统的两个有效模式——以及对 Claude Code 用户的启示 - 暗无天日 用 Org Babel 写 Literate 博文:扩展执行 + 定制导出 proced:Emacs 内置的进程查看器 - 暗无天日 从 proced 定制中学到的 Elisp 模式 读:让 Emacs proced 在 macOS 上显示 CPU 和内存 异步编程的函数着色税 - 暗无天日 链式调用的代价:JavaScript 和 Clojure 的共同教训 - 暗无天日 hyperfine:命令行基准测试工具 - 暗无天日 管道中的变量去哪了?——子 shell 作用域陷阱 - 暗无天日 开源包装器的信任陷阱:四个危险信号 - 暗无天日 程序员愿意为 AI 写文档,却不愿为同事写 - 暗无天日
TIL:流式处理的五个配置原则 - 暗无天日
lujun9972,Claude Code · 2026-06-03 · via 暗无天日

目录

  • 状态恢复不能依赖单节点
  • 写出频率和数据量要平衡
  • 分区数要匹配计算资源
  • 有状态操作必须设水位线控制状态增长
  • 三个监控指标上线前就搭好

DZone 一篇文章 中读到 Kuladeep Sandra 跑了五年 Kafka + Spark Structured Streaming 的生产经验。他处理的是保险理赔、制造业遥测、金融交易这类有 SLA 要求的系统。文章讲的是 Kafka/Spark 的具体配置,但这些原则的思路在其他流式处理框架中也能找到对应的做法。

状态恢复不能依赖单节点

流式处理框架用 checkpoint 记录两样东西:已提交的消费偏移量(offset)和有状态操作的中间状态。checkpoint 写在本地磁盘的话,进程重启能恢复,但节点挂了就没戏了。Spark 的做法是写到 HDFS、S3 等共享存储上;Kafka Streams 的做法是把状态变更写入 Kafka topic(changelog topic),靠 Kafka 的复制机制保证持久化。具体手段不同,但原则一样:状态恢复不能绑定在某一台机器上。

原文提到两次生产事故,一次是工程师把 checkpoint 目录当临时文件删了,一次是代码重构改了 query name 导致 checkpoint key 变化。两次都导致消费偏移量重置,需要从 Kafka 自身的 offset 存储中手工重建。都不致命,但都够折腾。所以别手动碰 checkpoint 目录,也别随便改影响 checkpoint key 的配置(比如 query name)。

写出频率和数据量要平衡

流式处理框架最终都要把数据写到下游存储。写得太频繁(比如每秒写一次),每次只攒了几条数据,产出大量小文件,下游查询和合并(compaction)压力大。写得太稀疏(比如攒 10 分钟再写),单次数据量大,处理时内存压力又上来了。

原文作者的经验法则是让每次写出的数据量落在 50-500MB 范围内。在 Spark Structured Streaming 里他用触发间隔控制这个节奏:制造业遥测管道用 30 秒触发,保险理赔管道用 2 分钟触发。其他框架的控制手段不同,但问题一样:得在延迟要求和下游存储效率之间找个平衡。

分区数要匹配计算资源

数据分片数和计算并行度是供需关系,哪边多了都是浪费。分片太少,计算资源吃不饱;分片太多,调度和上下文切换的开销反而拖慢处理。原文举的例子:Kafka 的每个分区在 Spark 中对应一个 task,如果 topic 有 200 个分区但集群只有 32 个核,200 个 task 在 32 个核上跑,上下文切换开销很大。反过来,如果下游处理比读取更吃 CPU,就得在读取后增加并行度——在 Spark 里是 repartition(),不过这会触发 shuffle,网络开销不小,只在确实需要时才用。

有状态操作必须设水位线控制状态增长

窗口聚合、流-流 join 这类有状态操作需要框架跨批次维护状态。不设水位线(watermark)的话,状态会无限增长,直到内存爆掉。

水位线阈值既是技术参数也是业务参数。设 10 分钟就是说「事件时间超过 10 分钟的迟到数据会被丢弃」。如果你的数据源本身就有 30 分钟延迟(IoT 设备批量上报、批结算场景),10 分钟水位线会把合法的迟到数据静默丢掉。因此,水位线就是你对迟到数据的容忍上限,设多长取决于你的数据源实际能迟到多久。

三个监控指标上线前就搭好

  1. 消费延迟(consumer lag),最重要的流式指标。延迟持续增长说明消费速度跟不上生产速度,SLA 快要违约了。
  2. 批处理耗时(batch duration),如果批处理耗时超过触发间隔,说明处理有瓶颈,作业已经跑不过来了。
  3. 状态存储大小(state store size),对有状态操作而言,状态持续增长就可能意味着内存泄漏,会 OOM。

这三个指标任何流式框架都能采集到,不需要特定云服务。关键是上线第一天就搭好,配上告警阈值(比如 consumer lag 连续 5 分钟增长就告警),别等出了第一次生产事故再补。