Recommended Free Tools
To write Kafka data to Delta Lake with exactly-once processing, use Spark Structured Streaming’s Delta sink with a durable, query-specific checkpoint, and keep the source offsets, checkpointed progress, and Delta commits aligned through recovery. Delta documents exactly-once processing at the Delta sink; that promise does not automatically cover arbitrary code, external side effects, or a Kafka output sink.
What exactly-once means in a Kafka-to-Delta pipeline
Exactly-once is a property of a defined processing path, not a blanket guarantee that every event exists only once everywhere. For this pipeline, the practical goal is that a failed micro-batch can be retried without losing committed Kafka input or applying the same output twice.
Delta Lake’s transaction log enables exactly-once processing for its Structured Streaming sink, including when other streams or batch queries use the table concurrently. Spark’s broader definition requires records to be received, transformed, and pushed downstream once; its output operations are at-least-once by default unless the output is idempotent or participates in a suitable transaction. See Delta Lake’s streaming documentation and the Apache Spark Streaming Programming Guide.
- Retry duplication is a repeated attempt to process the same Kafka offsets or micro-batch after failure. Checkpoint recovery and transactional or idempotent writes address this case.
- Duplicate source events are separate Kafka records representing the same real-world event. Offset-level exactly-once processing preserves both records; deduplicate by a genuine event identity if the application requires unique business events.
- External side effects such as API calls, database writes, or publishing to Kafka have their own delivery semantics. The Delta sink’s guarantee does not extend to them automatically.
How to set up the direct Structured Streaming path
The standard pattern is to read Kafka with Structured Streaming and write the resulting stream to Delta using writeStream.format("delta") and a persistent checkpointLocation. The checkpoint records streaming progress; Delta’s transaction log records committed table changes. Together they let the query recover across failures.
Do these 3 things before closing this tab:
1Clear out junk files and repair common Windows errors2Scan for outdated or missing drivers - takes under a minute3Repair Windows errors before they cause bigger problems#1 Best Overall
- Configure the Kafka source with the intended brokers and topic subscription, using the Kafka integration supported by the Spark version in your deployment.
- Apply transformations to the stream. If the business requirement is one row per real-world event, define an event identity and an appropriate deduplication rule; Kafka offsets alone do not establish business uniqueness.
- Write to Delta with the Delta streaming sink and a durable checkpoint path dedicated to this query. Keep the checkpoint accessible after a driver restart.
- Operate one active query per checkpoint. Do not start concurrent active queries from the same checkpoint location; Delta identifies this as a potential transaction-conflict scenario.
- Test recovery on production-equivalent storage. Confirm that a restart resumes from checkpointed progress and that storage supports Delta’s transactional requirements.
Use the Spark and Delta documentation for the exact configuration supported by the runtime versions you deploy. A checkpoint is operational state, not disposable temporary output: deleting or resetting it changes recovery behavior and may restart batch numbering.
When foreachBatch needs idempotency
foreachBatch gives a callback control over each micro-batch, but the callback may run again after a failure. Treat its effects as retryable rather than assuming the callback executes exactly once.
Idempotent Delta writes
Delta Lake documents idempotent table-write options for foreachBatch in Delta Lake 2.0.0 and later: use a stable txnAppId and a monotonically increasing txnVersion, such as the batch ID. Delta uses the application/version pair to recognize a repeated write and ignore it. See Delta Lake’s streaming documentation.
If you delete the checkpoint and start with a new one, assign a new application ID. The batch counter can begin again at zero; reusing the old ID could make new writes appear to be already-recorded transactions and cause them to be skipped.
What’s actually slowing this PC down?
Pick the symptom - the matching free tool is one click away.
Rank #3
MERGE and multiple targets
A MERGE inside foreachBatch must itself be idempotent: replaying a batch should converge on the intended table state rather than create duplicate effects. When a pipeline writes to several tables, separate streaming writes can provide better parallelization than serial writes in one callback; if a callback is necessary, make every target write retry-safe. Databricks discusses this guidance in its Lakeflow processing-guarantees documentation.
How Kafka offset commits fit in
For Structured Streaming’s integrated Kafka source and Delta sink, favor the query’s checkpoint-and-sink recovery path over custom offset handling. Validate any custom source or sink against the exact Spark version in use.
The Apache Spark Kafka integration guide describes three offset-storage strategies for its Spark Streaming integration: Spark checkpoints, Kafka’s offset commit API, or storing offsets in the same transaction as the results in a transactional data store. It warns that Spark output operations are at-least-once: checkpoints alone do not prevent repeated output, and committing offsets to Kafka is not itself atomic with the output. These distinctions matter especially when assessing legacy DStream examples or hand-managed offsets. See the Spark Streaming + Kafka Integration Guide.
Make the storage and recovery window part of the design
Delta’s ACID behavior depends on storage semantics: atomic visibility, mutual exclusion for final file creation, and consistent listing, or a suitable LogStore implementation. Local filesystem behavior may not provide the guarantees needed for concurrent transactional writes, so a successful local test does not validate the production storage layer. Check the requirements in Delta Lake’s storage configuration documentation.
Quick wins for a faster PC:
Clear out junk files and repair common Windows errorsFree Scan →Scan for outdated or missing drivers - takes under a minuteDriver Scan →Repair Windows errors before they cause bigger problemsFix Now →Best Value
Recovery also depends on source history and logs remaining available long enough to cover realistic outages and restart time. Delta warns that a streaming source falling behind cleaned transaction history may process only the latest available history and drop data. Databricks warns that a Delta stream beyond its data-file or log retention window may fail and require a full refresh. Set retention to cover the recovery window, and do not suppress missing-file errors with a setting that can silently return incomplete results. See Delta Lake’s streaming documentation and Databricks’ processing-guarantees guidance.
Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Choose a deployment approach based on operational needs
| Approach | Documented behavior | What to evaluate |
|---|---|---|
| Apache Spark Structured Streaming with Delta Lake | Open-source Spark/Delta path; the Delta sink uses transaction-log commits and checkpoints for its documented exactly-once processing guarantee. Delta Lake documentation | Runtime and library compatibility, storage and LogStore configuration, checkpoint operations, engineering ownership, and recovery procedures. |
| Databricks Lakeflow managed streaming tables | Databricks documents managed Kafka ingestion using Structured Streaming checkpoints and transactional Delta writes. Databricks documentation | Managed operations, deployment environment, governance and integration needs, recovery controls, and service cost. The cited documentation does not establish a comparable cost or performance result. |
The available official material does not establish a directly comparable Kafka-to-Delta throughput, latency, cost, or adoption benchmark. Choose based on your deployment requirements rather than assuming either approach is universally faster, cheaper, or safer.
Quick Recap
What to verify before calling the pipeline exactly-once
- The query uses a persistent checkpoint location that survives driver restarts and is not shared by concurrent active queries.
- The Delta sink commits are part of the same recovery design as the Kafka source offsets.
- Every callback or non-Delta side effect has a transaction, idempotency key, or downstream deduplication strategy.
- Reprocessing a batch cannot duplicate a business event when the event identity is known and uniqueness is required.
- Storage provides the visibility, exclusion, and listing behavior Delta requires, and retention covers realistic recovery time.
- Restart and replay behavior has been validated for the exact runtime, storage, and sink configuration in use.
Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.

