如何在Flink 1.4.0中查询可查询状态?文档代码及定义疑问
关于Flink Queryable State中key和jobId的定义疑问
我在查看Flink Queryable State相关代码时,发现这段代码的关键参数key和jobId没有明确的定义说明,而且在flink-queryable-state模块里也找不到可参考的测试用例,导致代码理解起来很困难:
QueryableStateClient client = new QueryableStateClient(tmHostname, proxyPort); // the state descriptor of the state to be fetched. ValueStateDescriptor<Tuple2<Long, Long>> descriptor = new ValueStateDescriptor<>("average", TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {}), Tuple2.of(0L, 0L)); CompletableFuture<ValueState<Tuple2<Long, Long>>> resultFuture = client.getKvState(jobId, "query-name", key, BasicTypeInfo.LONG_TYPE_INFO, descriptor); // now handle the returned value resultFuture.thenAccept(response -> { try { Tuple2<Long, Long> res = response.get(); } catch (Exception e) { e.printStackTrace(); } });
请问这段代码中的key和jobId是如何定义的?
关于jobId的定义
jobId是你要查询的Flink作业的唯一标识符,获取方式主要有三种:
- 作业提交时,命令行控制台会直接输出作业ID(格式是类似
7f6a5c4b-8d9e-1234-5678-90abcdef1234的UUID) - 登录Flink Web UI,每个运行中的作业卡片上都会显示对应的
Job ID - 如果是在代码里提交作业,
ExecutionEnvironment.execute()或StreamExecutionEnvironment.execute()会返回JobExecutionResult对象,调用getJobID()就能拿到:JobExecutionResult result = env.execute("My Queryable State Demo"); JobID jobId = result.getJobID();
关于key的定义
key就是你要查询的状态对应的键值对的键,必须和你的Flink作业中状态绑定的Key完全匹配:
- 比如你的作业是基于
KeyedStream处理数据,用keyBy()指定了某个字段作为Key(比如用户ID、订单ID),那这里的key就必须是该Key的具体取值 - 举个实际例子:如果你的作业按Long类型的用户ID做
keyBy,维护每个用户的平均统计状态,那这里的key就是你要查询的目标用户ID,比如12345L
补充:测试用例的替代方案
虽然flink-queryable-state模块自带的测试用例不多,但你可以自己搭个简单Demo验证:
- 先写一个带Queryable State的Flink作业:在
KeyedStream上调用asQueryableState("query-name", descriptor)把状态暴露出来 - 再用你贴的客户端代码,填入正确的
jobId和目标key,就能成功查询到对应状态了
内容的提问来源于stack exchange,提问作者Brutal_JL
相关产品推荐
相关产品推荐

