Flink keyBy 为什么要多做一次 MurmurHash?
在 Flink 中做 keyBy 的时候,内部根据 key 选择下游节点时,有这么一段逻辑,代码在 KeyGroupRangeAssignment.java#L75:
int computeKeyGroupForKeyHash(int keyHash, int maxParallelism) {
return MathUtils.murmurHash(keyHash) % maxParallelism;
}
简单来说,就是根据 key 的 hash 值,计算出 key 的分组。这里需要注意的是,并非直接简单地根据 hash 值取模,而是经过了一次 MurmurHash。
那为什么要这么做,解决了什么问题呢?
问题简单来说就是:keyBy 之后的分组要足够均匀,而用户定义的 keyHash 未必均匀。
举个最极端的例子:M % 10,但 M 的取值全都是 10 的倍数,比如 {10, 20, 30} 等。那无论 M 的取值有多分散和随机,M % 10 都等于 0,那么最终只有分组 0 里面有数据,将会发生严重的数据倾斜。
为什么 MurmurHash 能解决这个问题?