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

Flink流处理中如何获取Accumulator?求示例代码

Flink流处理场景下获取Accumulator结果的实现方案

流处理作业分为有界流(处理完数据自动终止)和无界流(持续运行)两种场景,对应不同的累加器结果获取方式,以下是具体实现:

1. 作业终止时获取(有界流/主动停止无界流)

如果是有界流作业,执行完毕后可直接通过execute()的返回值拿到JobExecutionResult;如果是无界流,主动cancel作业时也能获取包含累加器数据的结果。

示例代码:

import org.apache.flink.api.common.JobExecutionResult;
import org.apache.flink.api.common.accumulators.IntCounter;
import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class StreamAccumulatorDemo {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 定义带累加器的RichFunction
        RichMapFunction<String, String> lineCountFunc = new RichMapFunction<String, String>() {
            private final IntCounter lineCounter = new IntCounter();

            @Override
            public void open(Configuration parameters) throws Exception {
                super.open(parameters);
                // 注册累加器
                getRuntimeContext().addAccumulator("num-lines", lineCounter);
            }

            @Override
            public String map(String value) throws Exception {
                lineCounter.add(1);
                return value;
            }
        };

        // 构建有界流作业
        env.fromElements("a", "b", "c", "d")
           .map(lineCountFunc)
           .print();

        // 执行作业并获取结果
        JobExecutionResult result = env.execute("Stream Accumulator Job");
        // 读取累加器值
        int totalLines = result.getAccumulatorResult("num-lines");
        System.out.println("Total processed lines: " + totalLines);
    }
}

对于无界流(比如监听Kafka),可通过executeAsync()提交作业,主动cancel后获取结果:

// 提交无界流作业
JobClient jobClient = env.executeAsync("Unbounded Stream Job");

// 模拟运行一段时间后停止作业
Thread.sleep(10000);
JobExecutionResult result = jobClient.cancel().get();

// 获取累加器结果
int processedLines = result.getAccumulatorResult("num-lines");
System.out.println("Lines processed before cancel: " + processedLines);

2. 运行中实时获取(无界流监控)

持续运行的无界流无法通过终止作业的方式实时获取累加器,可通过Flink REST API查询当前值:

  • 接口地址:http://<flink-rest-host>:<port>/jobs/<job-id>/accumulators
  • 返回的JSON数据中包含所有已注册累加器的当前值。

示例代码(Java调用REST API):

import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.net.HttpURLConnection;
import java.net.URL;

public class AccumulatorRestQuery {
    public static void main(String[] args) throws Exception {
        String flinkRestAddress = "http://localhost:8081";
        String jobId = "your-job-id-here"; // 替换为实际作业ID

        URL apiUrl = new URL(flinkRestAddress + "/jobs/" + jobId + "/accumulators");
        HttpURLConnection conn = (HttpURLConnection) apiUrl.openConnection();
        conn.setRequestMethod("GET");

        BufferedReader reader = new BufferedReader(new InputStreamReader(conn.getInputStream()));
        String line;
        StringBuilder response = new StringBuilder();
        while ((line = reader.readLine()) != null) {
            response.append(line);
        }
        reader.close();

        // 可使用Jackson/Gson等JSON库解析结果,提取"num-lines"的值
        System.out.println("Accumulator status: " + response.toString());
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 06:57:38