Exactly-once end-to-end cần đủ ba điều kiện: nguồn đọc lại được (Kafka giữ message theo offset), job ghi nhớ đã xử lý tới đâu (checkpoint) và sink ghi idempotent hoặc có transaction.
Thiếu một trong ba thì chỉ đạt at-least-once.
from pyspark.sql import functions as F
raw = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", BROKERS)
.option("subscribe", "orders")
.option("startingOffsets", "earliest")
.option("maxOffsetsPerTrigger", 500_000)
.load())
orders = (raw.select(F.from_json(F.col("value").cast("string"), schema).alias("o"))
.select("o.*"))
(orders.writeStream.format("delta")
.option("checkpointLocation", "s3://lake/_checkpoints/bronze_orders")
.partitionBy("event_date")
.trigger(processingTime="1 minute")
.toTable("bronze.orders"))Cơ chế: mỗi micro-batch, Spark ghi khoảng offset sẽ xử lý vào checkpoint trước rồi mới ghi dữ liệu. Delta sink commit mỗi batch vào transaction log kèm batch id, nên khi job chết giữa chừng và chạy lại đúng batch đó, Delta nhận ra batch đã commit và bỏ qua. Offset được lưu trong checkpoint của Spark, không dựa vào offset commit của Kafka consumer group.
Khi cần upsert thay vì append: dùng foreachBatch với MERGE. Sink này mặc định chỉ at-least-once, nên phải tự làm idempotent: với Delta, truyền txnAppId và txnVersion (bằng batch id) khi ghi, hoặc viết MERGE theo khoá để chạy lại vẫn ra cùng kết quả.
Những chỗ vẫn sinh trùng dù pipeline đã exactly-once:
- Producer gửi lại cùng một sự kiện khi retry: dedup theo event_id, ví dụ dropDuplicatesWithinWatermark (Spark 3.5+) hoặc MERGE ở tầng silver.
- Xoá hoặc đổi đường dẫn checkpoint: job đọc lại từ startingOffsets.
Vận hành: trigger mỗi phút sinh nhiều file nhỏ nên cần lịch OPTIMIZE cho partition cũ; theo dõi consumer lag (offset mới nhất trên Kafka trừ offset đã xử lý); đặt retention của topic dài hơn thời gian job có thể dừng, nếu không dữ liệu hết hạn trước khi kịp đọc.
Lưu ý: Kafka sink của Structured Streaming chỉ đảm bảo at-least-once. Nếu đẩy kết quả ngược lại Kafka cho hệ khác đọc, consumer phía sau phải tự dedup.