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

SpringBoot API获取Flink结果问题:map函数内外finalList数据不一致

问题分析与解决方案

你遇到的问题本质是Flink的异步分布式执行模型和SpringBoot同步API模型的冲突,下面一步步拆解原因和解决办法:

为什么外部打印finalList为空?

  1. Flink作业是异步执行的:你在主方法中提交Flink作业后,代码会立即继续执行System.out.println("Final List: " + finalList);,此时Flink的map算子还没开始处理任何数据,自然finalList是空的。
  2. 线程/进程隔离:即使是本地运行Flink,map函数的代码也是在Flink的工作线程中执行的,和SpringBoot主线程是分开的,主线程不会等待Flink线程完成数据处理。如果是集群模式,Flink的TaskManager运行在单独的JVM进程中,你修改的finalList和主方法中的finalList根本不是同一个对象,完全无法共享。
  3. 流处理的无限性:Flink是处理无限数据流的,你的聚合会持续进行,不存在"处理完成"的时刻,所以主方法无法简单等待所有数据处理完再读取结果。

可行的解决方案

根据你的需求(SpringBoot API获取Flink实时聚合结果),推荐以下几种生产级方案:

Flink提供了Queryable State功能,允许外部应用直接查询Flink作业中的状态,这是最贴合流处理场景的方案。

步骤:

  1. 将聚合结果存入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());
                      }
                  }
              });
    
  2. 在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

  1. 添加Flink Redis依赖

    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-redis_2.12</artifactId>
        <version>1.15.0</version> <!-- 匹配你的Flink版本 -->
    </dependency>
    
  2. 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());
            }
        }
    }));
    
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 16:40:50