Will Spark Provide a Unified Exactly‑Once API for All Sinks?
0 reputation · 15 Oct 2020, 10:40 UTC
0 reputation · 15 Oct 2020, 10:40 UTC
Maintain exactly‑once semantics across all Structured Streaming sinks during retries, avoiding duplicate writes when a task restarts.
Apache Spark relies on checkpointing and write‑ahead logs to recover from failures. Exactly‑once guarantees are only available for sinks that implement idempotent writes (Kafka, Delta Lake, HDFS/Parquet). JDBC and custom sinks lack native idempotence, forcing developers to embed deduplication logic manually. The foreachBatch operator does not enforce any global guarantee, and the API exposes no single toggle to enable exactly‑once semantics across all sinks.
There is no clear roadmap indicating whether Spark will introduce a unified exactly‑once API that abstracts idempotent writes for arbitrary external systems. This decision impacts how developers design fault‑tolerant streaming pipelines.
29275 reputation · 15 Oct 2020, 19:50 UTC
No, Apache Spark is unlikely to provide a unified, global configuration toggle to enforce exactly-once semantics across all sinks. Because exactly-once delivery depends on the capabilities of the external storage system (the sink), Spark cannot abstract away the requirement for the sink to support either idempotent writes or atomic transactions.
Spark ensures exactly-once processing internally through checkpointing and write-ahead logs (WAL). However, exactly-once delivery is a distributed systems problem that requires coordination between the Spark driver and the external system.
foreachBatch Gap: This operator provides no native guarantees; the developer is responsible for implementing deduplication or transactional logic within the batch function.For a global "exactly-once" flag to work on a non-idempotent sink (like a standard JDBC connection), Spark would need to implement a Distributed Transaction Manager or a Two-Phase Commit (2PC) protocol across all possible external integrations. This would introduce massive performance overhead and require the external sinks to support XA transactions or similar protocols, which many NoSQL and legacy databases do not.
Since a unified API does not exist, you must implement fault tolerance based on the sink type:
foreachBatch to commit the batch only after the Spark offset is recorded.UPSERT (idempotent write) rather than an INSERT.Diagnostic Detail Needed: To provide a specific implementation pattern, please specify if your target sink supports UPSERT operations or atomic multi-row transactions.
Use comments to ask for clarification. Post a solution as an answer.
29,275 reputation · 15 Oct 2020, 19:07 UTC
Since the introduction of foreachBatch in Spark 2.0, the micro‑batch identifier (batchId) is passed to your function. You can leverage this monotonic long to make writes idempotent—for example, by including the batch ID in a composite primary key or by performing an UPSERT that filters on the batch ID. This does not alter Spark’s internal exactly‑once processing guarantee; it merely lets you build sink‑specific deduplication logic that works with the checkpoint‑driven replay of batches. Without such logic, a plain JDBC sink may still see duplicate rows on task retries, even though Spark’s processing remains exactly‑once.