Flink Keyed测试工具报空指针,求助正确依赖配置
Flink单元测试KeyedCoProcessFunction时Keyed State异常排查与解决
问题描述
测试KeyedCoProcessFunction实现类时,使用KeyedTwoInputStreamOperatorTestHarness初始化并打开测试工具时抛出以下异常:
java.lang.NullPointerException: Keyed state can only be used on a 'keyed stream', i.e., after a 'keyBy()' operation.
异常触发点在KeyedCoProcessFunction的open方法获取状态的代码处:
currentAssetNodeContext = getRuntimeContext().getState( new ValueStateDescriptor<>("currentAssetNodeContext", AssetNodeContext.class) ); // 此处抛出NullPointerException
项目中因在flink-streaming-java:1.16.1中找不到测试工具类,额外引入了flink-streaming-java_2.11:1.14.4的tests包,怀疑是依赖版本不一致导致问题。
问题根源
是的,跨版本依赖混用正是问题核心:
- Flink不同大版本的测试工具类(如
KeyedTwoInputStreamOperatorTestHarness)内部实现差异极大,1.14.4的测试工具无法适配1.16.1的KeyedCoProcessOperator与RuntimeContext初始化逻辑。 - 版本不匹配会导致
RuntimeContext未正确初始化keyed状态环境,调用getState()时触发空指针异常。 - 此外,Flink 1.16.x已不再支持Scala 2.11,引入
flink-streaming-java_2.11会进一步加剧依赖冲突。
正确的依赖配置(Flink 1.16.1)
Flink 1.16.1的KeyedTwoInputStreamOperatorTestHarness位于flink-streaming-java的tests classifier包中,替换旧版本依赖,添加以下测试依赖(根据项目Scala版本选择,1.16.x默认用Scala 2.12,若用2.13则修改为flink-streaming-java_2.13):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.12</artifactId> <version>1.16.1</version> <scope>test</scope> <classifier>tests</classifier> </dependency>
同时移除原有的flink-streaming-java_2.11:1.14.4依赖。
测试代码验证
你的测试代码逻辑本身无问题,确保OperationalDataAssetNodeIdKeySelector和AssetNodeParameterIdKeySelector能正确提取key即可。依赖修正后,重新运行测试,testHarness.open()将正确初始化keyed状态环境,避免空指针异常。
内容的提问来源于stack exchange,提问作者Peter C. Glade
相关产品推荐
相关产品推荐

