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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:27:07