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 能解决这个问题?
我们看下其源码,在 MathUtils.java#L137:
int murmurHash(int code) {
code *= 0xcc9e2d51;
code = Integer.rotateLeft(code, 15);
code *= 0x1b873593;
code = Integer.rotateLeft(code, 13);
code = code * 5 + 0xe6546b64;
code ^= 4;
code = bitMix(code);
if (code >= 0) {
return code;
} else if (code != Integer.MIN_VALUE) {
return -code;
} else {
return 0;
}
}
我们先思考,为什么会出现不均匀的问题?本质上在于 keyHash 的数值中一定隐藏了某种规律,恰好撞上分组规则,导致很多不同的 key 落进同一组中。比如上面提到的 key 都是 10 的倍数。
如果深入到 MurmurHash 的每一行实现,就非常复杂和难以解释了,我们这里只分析一下其原理。
MurmurHash 可以做到:
- 输出的 32 个 bit,每个 bit 可能受到输入的多个 bit 的影响。
- 反过来说,输入的每个 bit 的值,都可能影响输出的多个 bit。
从信息论的角度来看,MurmurHash 本质上做的是信息打散,即把原有每个 bit 的信息打散到多个 bit 上。
而源码中做乘法、移位、异或、旋转等操作,都是为了实现这个信息打散的目的。最终,某个输出位的值,就可能同时受输入中多个位的影响。这样,keyHash 里原本可能的隐藏规律就被打散了。
所以总结来说:Flink 在不信任用户 hashCode () 分布质量的前提下,增加一次确定性的位混合,来保证 keyBy 之后分组更均匀。