SpringBoot API获取Flink结果问题:map函数内外finalList数据不一致
问题分析与解决方案
你遇到的问题本质是Flink的异步分布式执行模型和SpringBoot同步API模型的冲突,下面一步步拆解原因和解决办法:
为什么外部打印finalList为空?
- Flink作业是异步执行的:你在主方法中提交Flink作业后,代码会立即继续执行
System.out.println("Final List: " + finalList);,此时Flink的map算子还没开始处理任何数据,自然finalList是空的。 - 线程/进程隔离:即使是本地运行Flink,map函数的代码也是在Flink的工作线程中执行的,和SpringBoot主线程是分开的,主线程不会等待Flink线程完成数据处理。如果是集群模式,Flink的TaskManager运行在单独的JVM进程中,你修改的
finalList和主方法中的finalList根本不是同一个对象,完全无法共享。 - 流处理的无限性:Flink是处理无限数据流的,你的聚合会持续进行,不存在"处理完成"的时刻,所以主方法无法简单等待所有数据处理完再读取结果。
可行的解决方案
根据你的需求(SpringBoot API获取Flink实时聚合结果),推荐以下几种生产级方案:
方案1:使用Flink Queryable State(查询状态)
Flink提供了Queryable State功能,允许外部应用直接查询Flink作业中的状态,这是最贴合流处理场景的方案。
步骤:
将聚合结果存入Flink的Managed State
不要用自定义的finalList,而是用Flink的MapState来存储聚合结果,这样状态会被Flink管理(容错、恢复)。// 定义状态描述符 MapStateDescriptor<String, Long> aggStateDesc = new MapStateDescriptor<>( "aggregated-results", BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.LONG_TYPE_INFO ); // 开启Queryable State aggStateDesc.setQueryable("queryable-agg-results"); // 在聚合算子后将结果存入MapState data_15min.keyBy(k -> "global") // 用一个固定key将所有数据分到同一个分区,适合全局聚合 .process(new ProcessFunction<Map<String, Long>, Void>() { private MapState<String, Long> aggState; @Override public void open(Configuration parameters) throws Exception { aggState = getRuntimeContext().getMapState(aggStateDesc); } @Override public void processElement(Map<String, Long> value, Context ctx, Collector<Void> out) throws Exception { // 更新状态 for (Map.Entry<String, Long> entry : value.entrySet()) { aggState.put(entry.getKey(), entry.getValue()); } } });在SpringBoot中查询Queryable State
使用Flink的Queryable State Client来查询最新的聚合结果:// 初始化Queryable State Client QueryableStateClient client = new QueryableStateClient("flink-jobmanager-host", 9069); // 查询状态 CompletableFuture<MapState<String, Long>> future = client.getKvState( JobID.fromHexString("your-job-id"), "queryable-agg-results", "global", // 和之前keyBy的固定key一致 BasicTypeInfo.STRING_TYPE_INFO, aggStateDesc ); // 等待查询结果(SpringBoot API中可以用异步处理) MapState<String, Long> resultState = future.get(); Map<String, Long> finalResult = new HashMap<>(); for (Map.Entry<String, Long> entry : resultState.entries()) { finalResult.put(entry.getKey(), entry.getValue()); } // 返回给API响应 return finalResult;
方案2:将Flink结果写入外部存储(推荐生产环境)
把Flink的聚合结果写入Redis、MySQL或其他KV存储,然后SpringBoot API直接从这些存储中读取数据,这种方式解耦了Flink和SpringBoot,扩展性更好。
示例:写入Redis
添加Flink Redis依赖
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-redis_2.12</artifactId> <version>1.15.0</version> <!-- 匹配你的Flink版本 --> </dependency>Flink中写入Redis
// 配置Redis连接 FlinkJedisPoolConfig jedisConfig = new FlinkJedisPoolConfig.Builder() .setHost("localhost") .setPort(6379) .build(); // 将聚合结果写入Redis data_15min.addSink(new RedisSink<>(jedisConfig, new RedisMapper<Map<String, Long>>() { @Override public RedisCommandDescription getCommandDescription() { return new RedisCommandDescription(RedisCommand.HSET, "flink-agg-results"); // 用Hash存储 } @Override public String getKeyFromData(Map<String, Long> data) { // 这里我们遍历每个entry写入,所以需要自定义处理 return null; // 这个方法会被覆盖,因为我们用批量处理 } @Override public String getValueFromData(Map<String, Long> data) { return null; } @Override public void open(Configuration parameters) throws Exception { super.open(parameters); } @Override public void write(Map<String, Long> data) throws Exception { // 遍历聚合结果,逐个写入Redis Hash for (Map.Entry<String, Long> entry : data.entrySet()) { getJedisCluster().hset("flink-agg-results", entry.getKey(), entry.getValue().toString()); } } }));SpringBoot中读取Redis
使用Spring Data Redis来读取Hash中的数据:@Autowired private StringRedisTemplate redisTemplate; @GetMapping("/agg-results") public Map<String, Long> getAggResults() { Map<String, String> redisResult = redisTemplate.opsForHash().entries("flink-agg-results"); // 转换为Long类型 return redisResult.entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, entry -> Long.parseLong(entry.getValue()) )); }
方案3:本地测试用异步阻塞(不推荐生产)
如果只是本地测试,可以用CountDownLatch来等待Flink处理一定数量的数据,但注意这只适合有限数据流的场景:
// 定义CountDownLatch,假设我们等待处理1条数据 CountDownLatch latch = new CountDownLatch(1); static Map<String, Long> finalList = new HashMap<>(); data_15min.map(new MapFunction<Map<String, Long>, Object>() { @Override public Object map(Map<String, Long> value) throws Exception { // 更新finalList逻辑 value.forEach((key, val) -> { finalList.put(key, val); }); latch.countDown(); // 处理完一条数据后递减计数器 return null; } }); // 提交Flink作业 env.execute(); // 等待计数器归零 latch.await(10, TimeUnit.SECONDS); // 设置超时时间,避免无限阻塞 // 现在可以获取finalList数据 System.out.println("Final List: " + finalList);
总结
- 生产环境优先选择方案2(写入外部存储),解耦性和扩展性最好;
- 如果需要直接和Flink交互,选择方案1(Queryable State);
- 方案3仅适合本地测试,不要用于生产。
内容的提问来源于stack exchange,提问作者Abhinav
相关产品推荐
相关产品推荐

