【流式数据处理】DataStream 与算子语义
拆解 Source/Transform/Sink 数据流图、rebalance/keyBy/broadcast 等 shuffle 策略、keyBy 到 KeyGroup 的映射,以及 ProcessFunction 与 TimerService 如何承载事件时间逻辑,并引入算子状态与键控状态的分工边界。
发布来自土法炼钢兴趣小组的知识、笔记、进展和应用。主题包括数据结构和算法、编程语言、网络安全、密码学等。
共 4 篇文章 · 返回首页
拆解 Source/Transform/Sink 数据流图、rebalance/keyBy/broadcast 等 shuffle 策略、keyBy 到 KeyGroup 的映射,以及 ProcessFunction 与 TimerService 如何承载事件时间逻辑,并引入算子状态与键控状态的分工边界。
拆解 Trino 的 partitioning scheme(HASH、BROADCAST、REPLICATE、ROUND_ROBIN)、LocalExchange 与 Remote Exchange、PartitionedOutput 数据路径,以及 skew 在 EXPLAIN ANALYZE 上的判读;对照 Spark shuffle 与 AQE 边界。
拆解 Spark 3.5+ Catalyst 的 Analyzed / Optimized / Physical 计划链、whole-stage codegen 与 shuffle 边界、AQE 的动态 coalesce/skew join/broadcast;并与 Trino 476 及 Iceberg V2 reader 下推能力对照。
闭合数据平台栈最后一块:从 SQL 解析与 Calcite 式优化,到 Volcano/向量化执行、Trino Coordinator/Worker 与 shuffle,再到 Iceberg connector 下推与生产排查。承接 lakehouse 第 18 章读湖视角,补全「谁在做 planning」的引擎内核层。