Spark Structured Streaming写入记录数统计异常排查求助
首先,你遇到的taskEnd.taskMetrics().outputMetrics().recordsWritten()始终为0的情况,不是你的流构造代码有问题,而是Spark Structured Streaming在文件类型Sink(比如你用的Parquet)下的一个已知问题——也就是你提交的SPARK-23288所指向的bug。这个bug的核心原因是:Structured Streaming的Append模式下,文件Sink的任务指标没有正确上报到TaskMetrics的OutputMetrics中,导致通过SparkListener捕获到的写入记录数始终为0,和实际写入情况不符。
为什么累加器显示“number of output rows = 2”?
你提到的累加器输出,应该是Spark内部中间算子的统计值(比如读取的输入行数),而非Sink实际写入的行数。仔细看你的代码,无效记录在catch块里没有被添加到输出列表,所以flatMap实际只输出1条有效记录,这个累加器的“2”和实际写入数无关。
可靠的读写记录数统计方法
针对Structured Streaming场景,推荐以下几种更准确的统计方式:
1. 利用StreamingQuery内置进度指标
Spark Structured Streaming的StreamingQuery对象提供了官方可靠的流处理进度API,能直接获取读写统计:
- 读取记录数:通过
query.lastProgress().numInputRows()获取,准确统计每个微批的输入行数。 - 写入记录数:根据Spark版本,通过
query.lastProgress().sinkMetrics().get("numOutputRows")(Spark 2.x)或query.lastProgress().numOutputRows()(Spark 3.x+)获取Sink实际写入的记录数。
示例代码片段:
StreamingQuery writeStream = ... // 你的流构造代码 writeStream.processAllAvailable(); // 获取最后一次微批的进度 StreamingQueryProgress progress = writeStream.lastProgress(); System.out.println("读取记录数: " + progress.numInputRows()); // 适配不同Spark版本的写入记录数获取 long writtenRows = Long.parseLong(progress.sinkMetrics().getOrDefault("numOutputRows", "0")); System.out.println("实际写入记录数: " + writtenRows); writeStream.stop();
2. 使用自定义累加器统计有效输出
如果需要在转换算子层面做细粒度计数,可以自定义累加器,在flatMap中对成功转换的记录进行统计:
// 定义长整型累加器 LongAccumulator validRecordAccumulator = session.sparkContext().longAccumulator("validRecordCount"); StreamingQuery writeStream = session .readStream() // ... 省略读取和配置代码 .as(Encoders.bean(IntegTestRecord.class)) .flatMap( ((FlatMapFunction<IntegTestRecord, IntegTestVendingRecord>) (u) -> { List<IntegTestVendingRecord> resultIterable = new ArrayList<>(); try { IntegTestVendingRecord result = transformer.convert(u); resultIterable.add(result); // 成功转换则累加计数 validRecordAccumulator.add(1); } catch (Throwable t) { System.err.println("Ooops"); t.printStackTrace(); } return resultIterable.iterator(); }), Encoders.bean(IntegTestVendingRecord.class)) // ... 省略写入配置代码 .start(); writeStream.processAllAvailable(); System.out.println("实际写入记录数(累加器统计): " + validRecordAccumulator.value()); writeStream.stop();
这种方式能精准统计经过转换后、实际要写入Sink的记录数,和最终写入结果完全一致。
3. 避免依赖SparkListener的TaskMetrics统计流写入
如前所述,TaskMetrics的OutputMetrics在Structured Streaming的文件Sink场景下存在统计不准确的bug,在对应的issue修复前,不建议用这种方式统计流的写入记录数。
总结一下:你遇到的是Spark的已知bug,而非代码错误;推荐使用StreamingQuery的进度指标或自定义累加器来实现可靠的读写记录数统计。
内容的提问来源于stack exchange,提问作者Yuriy Bondaruk

