Apache Flink中向Filter传递相同方法引用抛出ClassCastException问题
Flink中方法引用与Lambda调用filter的类型异常差异分析
在Flink 1.17.0 + Java 11环境下,出现了如下差异现象:
抛出ClassCastException的代码
final DataStreamSource<String> st1 = environment.fromCollection(List.of("asdf", "asfd")); final DataStreamSource<Integer> st2 = environment.fromCollection(List.of(1, 2, 3)); st1.filter(Objects::nonNull); st2.filter(Objects::nonNull);
运行时抛出异常:java.lang.ClassCastException: class java.lang.Integer cannot be cast to class java.lang.String
运行正常的代码
final DataStreamSource<String> st1 = environment.fromCollection(List.of("asdf", "asfd")); final DataStreamSource<Integer> st2 = environment.fromCollection(List.of(1, 2, 3)); st1.filter(s -> Objects.nonNull(s)); st2.filter(i -> Objects.nonNull(i));
结论:这是特性,而非Bug
核心原因是Java类型推断逻辑与Flink类型系统的交互特性:
- 使用
Objects::nonNull方法引用时,Java编译器会将其推断为Predicate<Object>类型(因为Objects.nonNull的参数是Object)。Flink基于这个类型生成序列化器时,会抹除流的具体元素类型信息,把所有元素当作Object处理。当运行时处理不同类型的流元素时,就会触发错误的跨类型强制转换。 - 使用Lambda表达式时,编译器会根据流的泛型上下文(
String/Integer),自动将Lambda推断为对应的Predicate<String>或Predicate<Integer>。Flink能正确识别每个流的具体元素类型,生成匹配的序列化器,因此不会出现类型转换异常。
内容的提问来源于stack exchange,提问作者rloeffel
相关产品推荐
相关产品推荐

