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

Spring Boot中调用Flink无界数据源聚合结果返回API问题

解决Spring Boot调用Flink无界流聚合结果的方案

方案1:利用Flink状态查询API直接获取

Flink内置状态查询能力,可让外部应用直接查询运行中作业的聚合状态,步骤如下:

    1. 在Flink作业中将聚合结果存入可查询的状态(KeyedState/OperatorState):
    public class AggregateProcessor extends RichMapFunction<RawData, AggResult> {
        private ValueState<AggResult> aggState;
    
        @Override
        public void open(Configuration params) throws Exception {
            ValueStateDescriptor<AggResult> stateDesc = new ValueStateDescriptor<>(
                "total-agg-result",
                TypeInformation.of(AggResult.class)
            );
            aggState = getRuntimeContext().getState(stateDesc);
        }
    
        @Override
        public AggResult map(RawData input) throws Exception {
            AggResult current = aggState.value() != null ? aggState.value() : new AggResult();
            // 执行累加、统计等聚合逻辑
            current.update(input);
            aggState.update(current);
            return current;
        }
    }
    
    1. 提交作业时开启检查点,保证状态持久化:
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    env.enableCheckpointing(3000); // 每3秒做一次检查点
    env.setRestartStrategy(RestartStrategies.fixedDelayRestart(2, 1000));
    
    1. Spring Boot端点中通过Flink客户端查询状态:
    @GetMapping("/allData")
    public AggResult getAggResult() throws Exception {
        ClusterClient<?> client = ClusterClientProvider.getClusterClient();
        JobID jobId = JobID.fromHexString("your-running-job-id");
        OperatorID operatorId = OperatorID.fromHexString("agg-processor-operator-id");
    
        CompletableFuture<AggResult> resultFuture = client.getKvState(
            jobId,
            operatorId,
            "total-agg-result",
            KeySelector.empty(),
            TypeInformation.of(AggResult.class)
        );
        return resultFuture.get();
    }
    

方案2:Flink将结果写入外部存储,Spring Boot查询存储返回

让Flink定期把聚合结果写入Redis、MySQL等外部存储,Spring Boot直接查询存储获取最新快照:

    1. 在Flink作业中添加Sink,定期输出聚合结果:
    dataStream.keyBy(RawData::getGroupKey)
              .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
              .aggregate(new CountAggregate())
              .addSink(new RedisSink<>(redisConfig, new AggResultRedisMapper()));
    
    1. Spring Boot端点直接查询外部存储:
    @GetMapping("/allData")
    public AggResult getAggResult() {
        String resultJson = redisTemplate.opsForValue().get("agg_total_key");
        return objectMapper.readValue(resultJson, AggResult.class);
    }
    
    这种方式实现简单,无需依赖Flink内部API,适合对实时性要求适中的场景。

如果用Flink SQL开发,可创建物化视图实时维护聚合结果,Spring Boot通过JDBC查询:

    1. 创建物化视图:
    CREATE MATERIALIZED VIEW user_behavior_agg
    AS SELECT user_id, COUNT(*) as action_count, MAX(ts) as last_action_time
    FROM user_behavior_stream
    GROUP BY user_id;
    
    1. Spring Boot通过Flink JDBC驱动查询:
    @GetMapping("/allData")
    public List<AggResult> getAggResult() throws SQLException {
        try (Connection conn = DriverManager.getConnection("jdbc:flink://flink-rest:8081")) {
            Statement stmt = conn.createStatement();
            ResultSet rs = stmt.executeQuery("SELECT * FROM user_behavior_agg");
            List<AggResult> results = new ArrayList<>();
            while (rs.next()) {
                results.add(new AggResult(
                    rs.getString("user_id"),
                    rs.getLong("action_count"),
                    rs.getTimestamp("last_action_time")
                ));
            }
            return results;
        }
    }
    

关键注意点

  • 无界流的聚合结果是动态更新的,每次API调用获取的都是当前时刻的最新快照。
  • 使用状态查询API时,需提前获取运行中作业的JobID和算子ID(可通过Flink UI或REST API查看)。
  • 外部存储方案需保证Flink的Exactly-Once语义,避免数据重复或丢失。

内容的提问来源于stack exchange,提问作者Abhinav

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 11:30:53