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

推荐订阅源

Recent Announcements
Recent Announcements
J
Java Code Geeks
雷峰网
雷峰网
Microsoft Security Blog
Microsoft Security Blog
博客园 - 【当耐特】
腾讯CDC
博客园 - 司徒正美
B
Blog RSS Feed
博客园 - 三生石上(FineUI控件)
I
InfoQ
N
Netflix TechBlog - Medium
L
LangChain Blog
博客园_首页
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
T
Tailwind CSS Blog
MyScale Blog
MyScale Blog
美团技术团队
The Cloudflare Blog
爱范儿
爱范儿
Stack Overflow Blog
Stack Overflow Blog
博客园 - 聂微东
H
Help Net Security
Martin Fowler
Martin Fowler
V
Visual Studio Blog

博客园 - 幽州散人

SQL优化案例 synchronized 为啥会让虚拟线程卸载失效 主线任务和使命 温州3日之旅 CompletableFuture + 异步Servlet:同步的姿势、异步的效果 Snowflake雪花算法与发号器的应用 windows Java 幽灵转发壳进程 也谈MySQL limit offset深翻页问题 SHOW SESSION STATUS 与 Handler 变量 MySQL 5.7 Profiles 执行性能分析 MySQL临时表与文件排序 MySQL 5.7 MRR多范围读优化 MySQL 5.7 ICP索引条件下推优化 MySQL符合索引与最左前缀原则 从I/O 的物理成本理解回表和索引覆盖 第4个电瓶以及补漆 什么是数字批发银行 Java虚拟线程(三)实现原理 Java虚拟线程(二) Java虚拟线程(一) 大模型推理层服务化架构 传奇调查员之路:理智值-99,但我必须听懂深渊的语言 js里调用智能合约读/写函数的方法的区别 字节与其16进制字符表示转换的bug Redisson分布式锁 交易心得 DexScreener接口初探 某安全软件跑飞了。。 Ed25519算法签名与验签的Java实现 Tendermint拜占庭容错引擎
Flink原理:并行度与keyGroup桶
幽州散人 · 2026-08-03 · via 博客园 - 幽州散人

数字银行风控场景:kafka -> Flink -> redis/hbase ,外部数据源先进kafka,Flink消费做流式聚合计算、出特征指标,然后存redis或hbase。 后边的各风控服务直接按需读指标、输入模型进行推理预测。

通常设置kafka Partition和parallelism相等

往kafka里写消息的时候,为了充分利用kafka分布式架构的优势、提高吞吐处理,可以按需设置分区进行并行写入,我们假设按照hash(企业ID)%3 进行3分区写入。
Flink侧一般也会按照同样数量设置parallelism,即source subTask数,这里也设置3,flink JM会在空闲的Slot中找3个Slot,也就是source subTask 0、1、2可能会分布在不同的TM

分区=并行度是为了 避免 Source 层的不均衡
Partition = Parallelism 解决的:
✅ 每个 Source Subtask 读 1 个分区,负载均衡
如果Partition ≠ Parallelism 会额外造成:
❌ 有的 Subtask 读 2 个分区(累死),有的空转(浪费)

这里很容易认为是这3个线程去分别处理对应的分区数据,之后走后续的流式计算,即按时间窗口聚合、按顺序计算聚合算子,然后sink。但是实际上不是这样。

两层分片和shuffle机制

前面说的source subTask相当于3个flink线程对kafka消费做了分片处理,parallelism除了是source subtask并行数之外,也代表window subtask数,即真正执行窗口聚合计算的线程。
但这里有第二次分片。murmurhash(企业ID) % 128, keyBy企业ID,把企业ID按照上述分片算法丢到128个keyGroup桶里,然后3个window subtask平分处理3组桶中的计算任务。

第一轮分片:Kafka 消费层

Producer(ENT001) → hash(ENT001) % 3 = 1 → Partition 1
Producer(ENT002) → hash(ENT002) % 3 = 0 → Partition 0
Producer(ENT004) → hash(ENT004) % 3 = 2 → Partition 2

          Partition 0      Partition 1      Partition 2
              │                 │                 │
              ▼                 ▼                 ▼
         Source Subtask 0  Source Subtask 1  Source Subtask 2
         (读到 ENT002)     (读到 ENT001)     (读到 ENT004)

第二轮分片:KeyBy 计算层


 

Source Subtask 0              Source Subtask 1              Source Subtask 2
  ENT002                        ENT001                        ENT004
    │                             │
    │  KeyBy(enterpriseId)        │                             │
    │  = MurmurHash % 128         │
    │                             │                             │
    ▼                             ▼                             ▼
  KeyGroup 11                   KeyGroup 82                   KeyGroup 55
  (归 Window Subtask 0)         (归 Window Subtask 2)         (归 Window Subtask 1)


                ┌──────────────────────────────────┐
                │         网络 Shuffle              │
                │                                  │
Source Subtask 0│  ENT002(KeyGroup 11)             │
                │     → 本地,因为 KeyGroup 11       │
                │       恰好归 Window Subtask 0     │
                │                                  │
Source Subtask 1│  ENT001(KeyGroup 82)             │
                │     → 跨线程,因为 KeyGroup 82     │
                │       归 Window Subtask 2        │
                │                                  │
Source Subtask 2│  ENT004(KeyGroup 55)             │
                │     → 跨线程,因为 KeyGroup 55     │
                │       归 Window Subtask 1        │
                └──────────────────────────────────┘

Window Subtask 0          Window Subtask 1          Window Subtask 2
KeyGroup [0..42]          KeyGroup [43..85]         KeyGroup [86..127]
   ENT002(本地)              ENT004(来自 S2)           ENT001(来自 S1)
   + 其他同范围企业          + 其他同范围企业          + 其他同范围企业
       ↓                         ↓
  窗口聚合 + Sink           窗口聚合 + Sink           窗口聚合 + Sink

两轮分片对比

┌────────┬──────────────────────────┬───────────────────────────────────┐
│        │     第一轮(Kafka)        │          第二轮(KeyBy)          │
├────────┼──────────────────────────┼───────────────────────────────────┤
│ 干嘛的  │ 把数据从 Kafka 读出来       │ 把数据按 key 路由给正确的计算线程     │
├────────┼──────────────────────────┼───────────────────────────────────┤
│ 分片数  │ 3(= 分区数)             │ 128(= maxParallelism)           │
├────────┼──────────────────────────┼───────────────────────────────────┤
│ 算法    │ hash(key) % 3            │ MurmurHash(key) % 128             │
├────────┼──────────────────────────┼───────────────────────────────────┤
│ 谁干    │ Source Subtask           │ KeyBy 算子                        │
├────────┼──────────────────────────┼───────────────────────────────────┤
│ 结果    │ 1 个 Subtask 读 1 个分区  │ 1 个 Subtask 扛 ~43 个 KeyGroup    │
└────────┴──────────────────────────┴───────────────────────────────────┘

Source Subtask ≠ Window Subtask。 虽然并行度都是 3,但它们是两套独立的线程。 中间 KeyBy 做了一次重新分配——企业 ID 在 Kafka 分区 1,不代表它的 KeyGroup 也归 Subtask1。 这个重新分配就是 "JM 下发任务"时的网络 shuffle。

之所以搞这么两层映射,是为了扩并行度方便,桶可以不变,线程重新分配自己处理的桶的范围就行了
128 个桶是固定底座,线程数可以变,桶的归属可以重新分配,但桶本身不动:

不搞两层映射:                        搞两层映射(KeyGroup):
key → hash % parallelism               key → hash % 128 → KeyGroup → Subtask

扩并行度: 3 → 6                       扩并行度: 3 → 6
hash(key) % 3 → hash(key) % 6         KeyGroup 桶不动
       ↑                                    ↑
  模数变了,每个 key 的分片结果全变了       只重新分桶的归属:
  → 状态全量迁移                         K0..K21 → S0    K22..K42 → S1
  → 必须停机重建                         K43..K63 → S2    K64..K85 → S3
                                         K86..K106 → S4   K107..K127 → S5

                                  状态跟着桶走,桶在哪个 Subtask 状态就在哪
                                  → 在线平滑扩缩,不停机

所以 KeyGroup 不是"搞复杂了",是一个变更隔离层——把"改并行度"这件对系统影响最大的操作,锁在一个可控的范围内。