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
相关产品推荐
相关产品推荐

