You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Flink 1.4.0中查询可查询状态?文档代码及定义疑问

我在查看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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 06:34:22