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

Spark重启后NATS Jetstream消息重复处理问题求助

解决NATS JetStream + Spark Streaming重复消费问题

问题场景

使用nats-spark-connector连接NATS JetStream,通过Spark Java代码消费并处理消息,代码如下:

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("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
            .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
            .getOrCreate();
    System.out.println("sparkSession : " + spark);
    
    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", "cons1")
            .option("nats.msg.ack.wait.secs", 120)
            .load();
    
    StreamingQuery query;
    try {
        query = df.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("delta")
                .option("path", "tmp/outputdelta")
                .start();
        query.awaitTermination();
    } catch (Exception e) {
        e.printStackTrace();
    } 
}

测试时通过NATS CLI推送10000条消息:

nats pub newsub --count=10000 "test #{{Count}}"

问题现象

Spark应用处理中途停止后重启,最终输出目录中的消息总数超过10000条(如测试中为10100条),出现重复。查看Delta Lake日志发现:应用停止前最后一个日志文件显示已处理到3100条,重启后第一个日志文件从3001条开始重新处理,导致3001-3100批次消息重复。

推测原因

使用Durable Consumer时,未发送确认(ACK)的消息会被NATS重新推送。应用关闭前的最后一个微批次尚未完成ACK,因此重启后该批次消息被重发并重复处理。

解决方案

1. 对齐Spark Checkpoint与NATS ACK时机

确认nats-spark-connector的ACK策略:确保只有当微批次数据成功写入Delta Lake后,才向NATS发送ACK。部分连接器默认消费到消息就立即ACK,需调整为批次级ACK——即整个微批次处理完成并持久化后,再批量确认该批次的所有消息。

2. 实现Spark优雅关闭

避免强制杀死应用,添加关闭钩子让应用完成当前批次的处理与ACK后再退出:

// 在创建query后添加关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    if (query != null) {
        query.stop();
        try {
            // 等待30秒让批次处理完成
            query.awaitTermination(30, TimeUnit.SECONDS);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
    spark.stop();
}));

3. 调整NATS ACK等待超时时间

当前nats.msg.ack.wait.secs设置为120秒,如果微批次处理时间超过该值,NATS会判定消息未处理完成并重发。根据实际处理耗时调整该参数,例如设置为300秒:

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

4. Delta Lake层面兜底去重

由于分布式系统中Exactly-Once语义难以绝对保证,在Delta Lake层通过唯一消息ID实现去重:

  • 首先从NATS消息中提取唯一标识(如消息ID)添加为DataFrame列;
  • 使用MERGE操作替代append,确保相同ID的消息不会重复插入:
DeltaTable deltaTable = DeltaTable.forPath(spark, "tmp/outputdelta");
StreamingQuery query = df
        // 提取NATS消息的唯一ID,具体字段需参考连接器文档
        .withColumn("msg_id", col("nats_msg_id"))
        .writeStream()
        .format("delta")
        .foreachBatch((batchDF, batchId) -> {
            deltaTable.as("target")
                    .merge(batchDF.as("source"), "target.msg_id = source.msg_id")
                    .whenNotMatchedInsertAll()
                    .execute();
        })
        .option("checkpointLocation", "tmp/checkpoint")
        .start();

内容的提问来源于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 01:27:02