Spark Streaming重启后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条消息,是否遗漏关键配置?
核心原因
- Spark Checkpoint与NATS ACK时序不匹配:Spark的checkpoint记录处理进度,但NATS持久化消费者的消息确认时机和Spark checkpoint保存不同步:
- 重启间隔过短:Spark还未将最新处理进度写入checkpoint就被终止,重启后从旧checkpoint位置重复处理,同时NATS未收到ACK会重新投递消息,导致重复计数。
- 重启间隔过长:超过
nats.msg.ack.wait.secs设置的90秒后,NATS会将未ACK的消息重新投递,但Spark从checkpoint读取的进度已经标记这些消息为已处理,导致这部分消息被跳过,总数不足。
- 文件输出的原子性问题:使用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();
验证步骤
- 重启前用
nats consumer info newstream newconsumer查看消费者未确认消息数,确认状态。 - 检查
tmp/checkpoint目录下的offset日志,验证Spark记录的处理进度。 - 调整配置后,测试不同重启间隔,确认最终处理总数为10000。
内容的提问来源于stack exchange,提问作者VGH

