Flink KeyBy中Lambda与双冒号方法引用的类型推断差异问题
Flink 1.10 keyBy Lambda类型推断失败根因
这个问题的本质是Flink 1.10版本类型提取器的实现局限,结合Java泛型擦除机制导致的,并非Java语言层面Lambda与方法引用存在语义差异。
- Flink流处理在构建StreamGraph阶段,必须明确获取KeyedStream的key类型,两个KeyedStream调用connect时会严格校验key类型一致性,类型不匹配就会抛出
java.lang.UnsupportedOperationException: Key types if input KeyedStreams don't match异常。 - 传入Lambda表达式
r -> r.getWorker_uniqe_id()作为KeySelector时,javac编译后的字节码中,Lambda对应的匿名实现类的泛型参数会被擦除到上界Object。Flink 1.10的TypeExtractor(类型提取组件)对Lambda表达式的字节码解析能力存在缺陷,无法穿透Lambda逻辑定位到实际调用的getWorker_uniqe_id()方法,因此只能读取到被擦除后的Object类型,将该流的key标记为GenericType<java.lang.Object>,最终触发类型不匹配错误。 - 传入双冒号形式的方法引用
WorkSign::getWorker_uniqe_id时,javac生成的字节码会保留方法引用指向的目标方法的明确签名元信息,Flink 1.10的TypeExtractor可以直接定位到WorkSign类下的getWorker_uniqe_id()方法,读取到方法定义的返回值为String类型,因此能正确识别key类型,后续connect校验时类型匹配,程序可正常运行。
1.10版本下使用Lambda编写KeySelector的兼容方案,可显式声明返回类型绕过类型提取缺陷:
keyBy(r -> r.getWorker_uniqe_id()).returns(String.class)
内容的提问来源于stack exchange,提问作者kursk.ye
相关产品推荐
相关产品推荐

