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

