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

Spark Streaming重启后NATS JetStream消息处理数量不准确问题

问题:Spark集成NATS JetStream断点续传后消息计数不准确

场景与问题

我使用以下Spark Java代码从NATS JetStream拉取消息:

private static void sparkNatsTester() {
         
            SparkSession spark = SparkSession.builder()
                    .appName("spark-with-nats")
                    .master("local")
                      .config("spark.jars",
                      "libs/nats-spark-connector-balanced_2.12-1.1.4.jar,"+"libs/jnats-2.17.1.jar"
                      )
                     .config("spark.sql.streaming.checkpointLocation","tmp/checkpoint")
                    .config("fs.s3a.access.key", "minioadmin")
                    .config("fs.s3a.secret.key", "minioadmin")
                    .config("fs.s3a.endpoint", "http://127.0.0.1:9000")
                    .config("fs.s3a.connection.ssl.enabled", "true")
                    .config("fs.s3a.path.style.access", "true")
                    .config("fs.s3a.attempts.maximum", "1")
                    .config("fs.s3a.connection.establish.timeout", "5000")
                    .config("fs.s3a.connection.timeout", "10000")
                      .getOrCreate();
            Dataset<Row> df = spark.readStream()
                    .format("nats")
                    .option("nats.host", "localhost")
                    .option("nats.port", 4222)
                    .option("nats.stream.name", "newstream")
                    .option("nats.stream.subjects", "newsub")
                    .option("nats.durable.name", "newconsumer")
                    // wait 90 seconds for an ack before resending a message
                    .option("nats.msg.ack.wait.secs", 90)
                    .load();
            System.out.println("Successfully read nats stream !");
            df.createOrReplaceTempView("natsmessages");
            Dataset<Row> filteredDf = spark.sql("select * from natsmessages");
            
            StreamingQuery query;
            try {
                query = filteredDf.withColumn("date_only", from_unixtime(unix_timestamp(col("dateTime"),  "MM/dd/yyyy - HH:mm:ss Z"), "MM/dd/yyyy"))
                        .writeStream()
                          .outputMode("append")
                          .partitionBy("date_only")
                          .format("parquet")
                          .option("path", "tmp/newoutputtest")
                          .start();
                query.awaitTermination();
            } catch (Exception e) {
                e.printStackTrace();
            } 
        } 

测试流程:通过NATS CLI执行nats pub newsub --count=10000 "test #{{Count}}"推送10000条消息,当应用处理2000条后关闭Eclipse中的应用并重启。已配置spark.sql.streaming.checkpointLocation,预期重启后从断点继续处理剩余消息,但最终总处理数不符合预期:

  • 重启间隔较长时,总处理数不足10000
  • 重启间隔极短时,总处理数超过10000

注:newconsumer是NATS中创建的持久化消费者。需要解决如何确保最终恰好处理10000条消息,是否遗漏关键配置?


核心原因

  1. Spark Checkpoint与NATS ACK时序不匹配:Spark的checkpoint记录处理进度,但NATS持久化消费者的消息确认时机和Spark checkpoint保存不同步:
    • 重启间隔过短:Spark还未将最新处理进度写入checkpoint就被终止,重启后从旧checkpoint位置重复处理,同时NATS未收到ACK会重新投递消息,导致重复计数。
    • 重启间隔过长:超过nats.msg.ack.wait.secs设置的90秒后,NATS会将未ACK的消息重新投递,但Spark从checkpoint读取的进度已经标记这些消息为已处理,导致这部分消息被跳过,总数不足。
  2. 文件输出的原子性问题:使用Parquet输出时,如果没有配置正确的提交器,可能出现部分写入的文件,重启后重复写入导致计数异常。

解决方案与配置调整

1. 对齐NATS ACK与Spark Checkpoint时机

NATS Spark连接器支持通过配置将ACK时机绑定到Spark Checkpoint完成后,添加以下配置到流读取部分:

.option("nats.ack.mode", "checkpoint")

该配置表示只有当Spark成功完成checkpoint保存后,才会向NATS发送消息确认,避免进度不一致。

2. 调整NATS ACK超时时间

将nats.msg.ack.wait.secs设置为大于Spark的checkpoint间隔加上单批消息处理时间,避免NATS在Spark未完成checkpoint就重新投递消息,例如调整为120秒:

.option("nats.msg.ack.wait.secs", 120)

3. 优化Spark Checkpoint配置

缩短checkpoint间隔,确保处理进度及时保存,在SparkSession创建时添加:

.config("spark.sql.streaming.checkpointInterval", "10s")

同时保留足够的checkpoint批次历史:

.config("spark.sql.streaming.minBatchesToRetain", "20")

4. 配置文件输出的原子提交器

针对S3a存储的Parquet输出,配置魔法提交器确保文件写入的原子性,避免部分写入导致重复计数,在SparkSession创建时添加:

.config("fs.s3a.committers.impl", "org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory")
.config("fs.s3a.committer.name", "magic")
.config("fs.s3a.committer.magic.enabled", "true")

5. 使用ForeachBatch实现精确控制(可选)

如果需要更精细的控制,可以用foreachBatch手动处理数据写入和NATS ACK:

query = filteredDf.withColumn("date_only", from_unixtime(unix_timestamp(col("dateTime"),  "MM/dd/yyyy - HH:mm:ss Z"), "MM/dd/yyyy"))
        .writeStream()
        .outputMode("append")
        .foreachBatch((batchDF, batchId) -> {
            // 原子写入数据
            batchDF.write()
                   .mode("append")
                   .partitionBy("date_only")
                   .parquet("tmp/newoutputtest");
            // 如果连接器提供手动ACK API,在此调用确认当前批次消息
            // 例如:batchDF.sparkSession().conf().set("nats.manual.ack", "true");
        })
        .start();

验证步骤

  1. 重启前用nats consumer info newstream newconsumer查看消费者未确认消息数,确认状态。
  2. 检查tmp/checkpoint目录下的offset日志,验证Spark记录的处理进度。
  3. 调整配置后,测试不同重启间隔,确认最终处理总数为10000。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 21:47:08