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 可以做到:

  1. 输出的 32 个 bit,每个 bit 可能受到输入的多个 bit 的影响。
  2. 反过来说,输入的每个 bit 的值,都可能影响输出的多个 bit。

从信息论的角度来看,MurmurHash 本质上做的是信息打散,即把原有每个 bit 的信息打散到多个 bit 上。

而源码中做乘法、移位、异或、旋转等操作,都是为了实现这个信息打散的目的。最终,某个输出位的值,就可能同时受输入中多个位的影响。这样,keyHash 里原本可能的隐藏规律就被打散了。

所以总结来说:Flink 在不信任用户 hashCode () 分布质量的前提下,增加一次确定性的位混合,来保证 keyBy 之后分组更均匀。