Flink v1.4.0可查询状态(QS)遇空指针异常求助
排查Flink 1.4.0可查询状态(QS)升级后失效的问题
从你的描述和错误栈来看,主要有两个关键问题导致升级后任务失效,我们逐一分析并给出修复方案:
1. QueryableStateClient API的参数类型变化
Flink 1.4.0对QueryableStateClient.getKvState()方法做了API变更:第一个参数从字符串类型的jobId改为了JobID对象,而你的代码依然传递字符串类型的jobId,这会导致运行时参数不匹配,进而引发序列化环节的空指针异常(NPE)。
修复方案:
将字符串jobId转换为JobID对象后再传入:
QueryableStateClient client = new QueryableStateClient("localhost", 9069); // 注意端口也需要修改,见下文 ValueStateDescriptor<DataBucket> descriptor = new ValueStateDescriptor<>( "QSName", TypeInformation.of(new TypeHint<DataBucket>() {})); String jobIdStr = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"; JobID jobId = JobID.fromHexString(jobIdStr); // 转换为JobID对象 String key = "2017-01-06"; CompletableFuture<ValueState<DataBucket>> resultFuture = client.getKvState( jobId, "QSName", key, BasicTypeInfo.STRING_TYPE_INFO, descriptor);
2. Queryable State连接端口变更
Flink 1.3.2中Queryable State客户端直接连接JobManager的6123端口,但1.4.0引入了Queryable State Proxy组件,客户端需要连接TaskManager上的Proxy端口,默认是9069(可通过flink-conf.yaml中的queryable-state.proxy.port配置自定义端口)。
如果你的客户端仍然连接6123端口,请求根本无法到达QS服务端,这也是任务失效的潜在原因。
修复方案:
修改客户端的连接端口为9069(或你配置的自定义端口):
QueryableStateClient client = new QueryableStateClient("localhost", 9069);
额外验证点
- 确保
DataBucket类在客户端和服务端的类路径完全一致(类版本、序列化方式统一),避免序列化/反序列化失败。如果DataBucket没有实现Serializable,建议添加该接口,或为其自定义FlinkTypeSerializer。 - 再次确认TaskManager已加载
flink-queryable-state-runtime_2.11-1.4.0.jar,可以查看TaskManager的日志验证组件是否成功初始化。
按照以上步骤修改后,应该能解决你遇到的NPE和QS请求无法到达的问题。
内容的提问来源于stack exchange,提问作者Christos Hadjinikolis
相关产品推荐
相关产品推荐

