DEV Community

Cover image for Spark Structured Streaming Can’t Save You From Bad Architecture
Aniket Abhishek Soni
Aniket Abhishek Soni

Posted on

Spark Structured Streaming Can’t Save You From Bad Architecture

Roughly 80% of the "exactly-once" guarantees I see implemented in production financial pipelines are actually "at-least-once" disguised by a fragile post-processing cleanup script. We obsess over the Spark checkpointing mechanism while ignoring the fact that our sink is a leaky bucket.

Why I chose this topic: I’ve spent the last six years cleaning up "exactly-once" messes that caused multi-million dollar reconciliation failures in healthcare billing systems. If you don't understand the physical limitations of your sink, your streaming job is just a very expensive way to generate duplicate records.

Most of my peers argue that Spark’s checkpointLocation is the gold standard for fault tolerance. They treat it like a holy relic. The reality is that Spark’s exactly-once guarantee is strictly confined to the boundary between the source offset management and the state store. Once your data hits an external database or a file system that doesn't participate in a distributed transaction with Spark, your "exactly-once" guarantee is effectively vaporware.

Why the common approach falls short

The industry standard is to use foreachBatch to sink data into a relational database or a cloud object store. Developers assume that because Spark retries the batch upon failure, they are protected against duplicates. This is a dangerous fallacy.

If your Spark executor writes 10,000 rows to a Postgres table and crashes at row 9,999, the driver will restart the task. If you aren’t using idempotent writes—where the sink itself knows how to deduplicate based on a business key—you now have a partial write. If your sink logic involves an INSERT statement without an ON CONFLICT DO UPDATE or a similar upsert mechanism, you have just injected garbage into your downstream analytical layer.

Take a look at a common, broken pattern I see in Spark 3.3.x:

streamingDF.writeStream
  .foreachBatch { (batchDF: DataFrame, batchId: Long) =>
    batchDF.write.format("jdbc").mode("append").save(...)
  }
  .start()
Enter fullscreen mode Exit fullscreen mode

This code is a lie. If the save() operation fails midway, Spark reruns the batch. If the database processed 5,000 rows before the crash, those 5,000 rows remain. The retry will attempt to insert all 10,000 rows again. You have just doubled your data. Unless your schema has a strict unique constraint and you’ve configured your JDBC driver to handle the resulting exception, you’ve corrupted your downstream data.

Photo by CHUTTERSNAP on Unsplash
Photo by CHUTTERSNAP on Unsplash

The illusion of checkpointing

Spark’s checkpointing is brilliant for maintaining the state of your streaming query. It stores the metadata of processed offsets in HDFS or S3. But people confuse "I know which offsets I’ve processed" with "I have committed the data to the sink."

When a worker node dies, Spark rolls back the state to the last successful offset recorded in the checkpoint directory. It then re-executes the batch. This is where the physical reality of your sink takes over. If you are writing to S3 using the default file sink, Spark uses a commit protocol that involves renaming temporary files. If that rename operation isn't atomic—and on many S3-compatible object stores, it isn't—you are left with "ghost" files that weren't cleaned up by the failed task but aren't technically part of the final manifest.

In my experience, relying on checkpointLocation without a high-performance, atomic storage layer is like trying to build a skyscraper on a swamp. You aren't doing exactly-once; you are doing "at-least-once plus a prayer that the infrastructure handles the retry gracefully."

The sink is the bottleneck

True exactly-once in a streaming system requires a two-phase commit protocol or an idempotent sink. Spark Structured Streaming implements this only when both the source and the sink support it natively. For example, if you are reading from Kafka and writing to Kafka, Spark can participate in the transaction.

But when you write to a sink like Delta Lake, you are relying on the Delta protocol to handle the atomicity. Delta Lake is excellent—it uses a transaction log to ensure that a batch is either fully committed or not at all. If you are using Delta Lake 2.0+, the writeStream API will correctly ignore duplicate batches because it checks the batchId against the transaction log before committing.

The problem arises when engineers try to shoehorn this into non-transactional sinks. I’ve seen teams try to implement "exactly-once" against Redis or a legacy REST API. They write custom logic inside foreachBatch to track batchId in a side-table. This is a recipe for a race condition. If your side-table update and your sink write aren't in the same transaction, you are fundamentally broken. You’ve just shifted the failure mode from "duplicate data" to "stale state."

Photo by Lukas Tennie on Unsplash
Photo by Lukas Tennie on Unsplash

The objections (and my answers)

"But the documentation says Spark guarantees exactly-once!"

The documentation says Spark guarantees exactly-once processing semantics. That is a technical term of art, not a business guarantee. It means that the state of the Spark job will be consistent with the offsets of the source. It explicitly avoids making promises about the side effects of your code. If your side effect is "inserting a row into a database," Spark isn't responsible for that database's internal state.

"Can't I just use a deduplication buffer?"

Sure, you can use dropDuplicates() on your streaming DataFrame. But dropDuplicates() is stateful. It has to keep a history of keys in memory or in the state store. If your window is large or your data volume is high, you will OOM (Out of Memory) your executors or blow up your RocksDB state store. You’re trading a duplication problem for a scaling problem.

"What about the newer Spark versions with better fault tolerance?"

Spark 3.5+ has improved the way it handles state store snapshots and recovery, but the sink problem remains. No matter how much you optimize the Spark engine, if you point it at a sink that doesn't support transactions, you are living in a world of at-least-once. The only way to achieve truly robust pipelines is to embrace idempotent sinks and design your schema so that re-processing the same batch is a no-op, not a disaster.

Conclusion

Stop chasing the ghost of "exactly-once" in your application code. You cannot force a non-transactional system to become transactional just by wrapping it in a Spark job.

Instead, build for idempotency. Use unique business keys in your database. Use Delta Lake or Hudi if you need transactional file storage. If you are writing to an API, ensure your request includes a client-side idempotency token that the API respects.

Spark Structured Streaming is a phenomenal tool for stateful processing and offset management. But it is not a magic wand that absolves you from the responsibility of handling data integrity at the sink. If you assume the framework handles everything, you’ll spend your weekends doing manual data reconciliation. Don’t trust the abstraction—trust the physical limitations of your infrastructure.

Cover photo by Tyler on Unsplash.

Top comments (0)