A Structured Streaming query is only as reliable as its source. The engine's fault-tolerance story, that a restarted query neither loses nor double-processes input, rests on a contract every source must honour: it must describe its input as offsets, and it must return the same data whenever it is asked for the same offset range. Sources that cannot do that are labelled not fault-tolerant in Spark's own documentation, and they exist only for testing.
This article explains that contract, then walks the built-in sources in Spark 4.x: the file source in depth, the test sources (rate, rate-micro-batch and socket), table sources, and where Kafka fits. It covers admission control, a worked backlog-recovery example, the Python Data Source API for writing your own source, and the failure modes that silently skip or re-read input. Kafka gets a short section here because Spark Streaming Kafka Source covers it fully.
The source contract
Every micro-batch begins with the driver asking each source for its latest available offset. An offset is whatever the source uses to name a position: a Kafka partition-to-offset map, an index into a log of seen files, a counter. Before reading anything, the driver writes the chosen range for batch N into offsets/N in the checkpoint directory. Executors then read exactly that range and the sink writes the result tagged with batch ID N. Only after the sink succeeds does the driver write commits/N and tell the source it may release data up to the end offset.
After a crash, Spark looks for the highest batch with an offset entry but no commit entry and re-runs it with the logged range. Two properties make that safe. The source must be replayable: asked for the same range, it returns the same records. And the sink must be idempotent per batch ID, which is the subject of Spark Streaming Sinks. Together they give end-to-end exactly-once results. Break either and you get duplicates or gaps with no error.
The built-in sources
| Source | Offset is | Fault-tolerant | Use it for |
|---|---|---|---|
File (text, csv, json, parquet, orc) | Index into a log of files already seen | Yes | Landing zones, exports, object-store drops |
| Kafka | Per-partition offsets | Yes | Event streams |
Table (readStream.table) | Table version or snapshot | Depends on the table format | Delta, Iceberg and similar as a change feed |
| Rate | Elapsed time and rows per second | Yes | Load tests |
Rate per micro-batch (rate-micro-batch) | Batch count | Yes | Deterministic tests |
| Socket | None that survives restart | No | Demos only |
Notice that "fault-tolerant" is about the source's ability to replay, not about whether your query is correct. A file source is fault-tolerant, yet a producer that overwrites a file in place can still make a replay read different bytes. The contract binds the producer too.
The file source in depth
The file source watches a directory (glob patterns are allowed; comma-separated lists of paths are not). On each trigger it lists the directory, subtracts the files it has already recorded in its metadata log under sources/0 in the checkpoint, and makes the remainder available as the next batch. A file is processed once, identified by path, so the right way to deliver data is to write a complete file and then make it visible, for example by writing to a temporary name and renaming on HDFS, or by uploading whole objects to an object store where an object appears only when fully written.
File sources require an explicit schema by default. Spark could infer one from the first files, but a later restart could then infer a different one; spark.sql.streaming.schemaInference re-enables inference for ad-hoc use. Spark recurses into partition directories such as date=2026-10-04/, but fills the partition column only if it appears in your schema, so add it there. The partitioning scheme must exist when the query starts and stay static: new values are fine, a new partition column is not.
from pyspark.sql.types import StructType, StructField, StringType, LongType, TimestampType
schema = StructType([
StructField("event_id", StringType(), False),
StructField("user_id", StringType()),
StructField("amount_cents", LongType()),
StructField("event_time", TimestampType()),
])
events = (spark.readStream
.schema(schema)
.option("maxFilesPerTrigger", 500) # admission control
.option("cleanSource", "archive") # move processed files away
.option("sourceArchiveDir", "s3a://lake-archive/landing/events")
.option("mode", "PERMISSIVE")
.json("s3a://lake-landing/events/"))The options worth knowing, with defaults from the Spark 4.2 guide:
maxFilesPerTrigger/maxBytesPerTrigger: cap new input per batch (no cap by default). Only one may be set, and a batch always takes at least one file so a single oversized file cannot stall the query.latestFirst: process newest files first when there is a backlog (default false). Useful when fresh data matters more than order.maxFileAge: ignore files older than this, measured against the newest file seen, not the wall clock (default one week). Ignored whenlatestFirstis combined with a per-trigger cap.fileNameOnly: identify files by name alone, so the same name under different paths or URI schemes counts as one file (default false).maxCachedFilesanddiscardCachedInputRatio: how many listed-but-unprocessed files are kept between batches (10,000 and 0.2) so that rate-limited queries do not relist every trigger.cleanSource:off,archiveordeletefor files whose batch has committed. Archiving needssourceArchiveDir, which must not match the source glob at the same depth; files keep their relative path under it.
Kafka, table and test sources
Kafka is the most common production source. Its offsets are per partition, replay comes from Kafka's retention, and maxOffsetsPerTrigger provides admission control. The catch is that replay only works while the data is still retained; a query stopped longer than the retention period resumes into a gap.
Table sources, via spark.readStream.table(...) or the format's own reader, turn a table format's commit history into a stream. They inherit the format's rules about what happens when a commit rewrites or deletes data, which usually needs an explicit option to skip or the stream fails.
Rate generates timestamp and value columns at rowsPerSecond, optionally ramping up over rampUpTime, across numPartitions partitions; use it to find how many rows per second a pipeline sustains. Rate per micro-batch emits rowsPerBatch rows each batch with generated time advancing by advanceMillisPerBatch from startTimestamp, so output is independent of wall-clock speed, which makes it the right tool for testing watermark and window logic. Socket reads lines from a TCP server, has no replayable offsets, and should never leave a notebook.
Admission control and triggers
Without limits, the first batch after an outage tries to read the entire backlog, which can exceed executor memory or produce one batch so long that downstream freshness alarms fire for an hour. Per-trigger caps make the source report a smaller end offset, so the backlog drains as a sequence of normal-sized batches.
Triggers interact with this. Trigger.AvailableNow processes everything available at start, in multiple batches that respect the source's caps, then stops; it is the right way to run a streaming query as a scheduled job. The older Trigger.Once processed everything in a single batch and ignored the caps, and is deprecated in its favour. Trigger choice is covered in Spark Streaming Triggers.
Worked example: draining a 30-hour backlog
A team lands JSON event files in an object-store prefix: about 40,000 files an hour, averaging 2 MB. The streaming job that reads them was down for 30 hours after a bad deploy. At restart there are 1.2 million unread files, about 2.4 TB. The cluster comfortably processes 100 GB per batch in roughly four minutes.
Uncapped, batch one would try to plan and read 2.4 TB at once. Capped with maxBytesPerTrigger at 100 GB, the query drains about 1.5 TB an hour while new data arrives at 80 GB an hour, a net 1.42 TB an hour. The 2.4 TB backlog is gone in about 1.7 hours, roughly 26 batches, each one an ordinary size that the cluster has already proven it can handle. If dashboards need fresh numbers immediately, latestFirst would show recent data first, but it disables maxFileAge filtering and reorders events, so any event-time logic must already tolerate disorder via a watermark.
One trap in this example: the operators restore a day of files from a backup bucket with a tool that preserves original modification times. Those files are more than a week older than the newest file in the directory, so maxFileAge silently skips them. Either copy them with fresh timestamps or run a separate bounded backfill query with its own checkpoint.
Observing sources in production
Every source reports per-batch progress, and that is where source problems show up first. Each entry in query.lastProgress["sources"] carries the start and end offsets of the batch, the number of input rows, and input versus processed rows per second. Input rate persistently above processed rate means a growing backlog; zero input rows for a source that should be busy means the producer, the path or the permissions changed.
for src in query.lastProgress["sources"]:
lag_signal = src["inputRowsPerSecond"] - src["processedRowsPerSecond"]
emit_metric("stream.input_rows", src["numInputRows"], tags={"source": src["description"]})
emit_metric("stream.rate_gap", lag_signal, tags={"source": src["description"]})For the file source, also watch the checkpoint's seen-files log. It grows with every file ever processed and is compacted periodically, but a query that has consumed hundreds of millions of small files carries that history into every restart.
Writing your own source in Python
When no built-in source fits, Spark 4.0 added a Python Data Source API with a streaming reader. You implement the same contract described above: an initial offset, a latest offset, partition planning for a range, a reader for each partition, and a commit hook. Offsets are JSON-serializable dicts that Spark writes to the offset log for you.
from pyspark.sql.datasource import DataSource, DataSourceStreamReader, InputPartition
class RangePartition(InputPartition):
def __init__(self, start, end):
self.start, self.end = start, end
class LedgerStreamReader(DataSourceStreamReader):
def __init__(self, schema, options):
self.api = options["endpoint"]
def initialOffset(self):
return {"seq": 0}
def latestOffset(self):
return {"seq": fetch_head_sequence(self.api)} # must be durable, ordered positions
def partitions(self, start, end):
step = 10_000
return [RangePartition(s, min(s + step, end["seq"]))
for s in range(start["seq"], end["seq"], step)]
def read(self, partition):
for rec in fetch_range(self.api, partition.start, partition.end): # same range -> same rows
yield (rec["seq"], rec["account"], rec["amount"])
def commit(self, end):
pass # e.g. allow upstream to trim before end["seq"]
class LedgerSource(DataSource):
@classmethod
def name(cls):
return "ledger"
def schema(self):
return "seq long, account string, amount long"
def streamReader(self, schema):
return LedgerStreamReader(schema, self.options)
spark.dataSource.register(LedgerSource)
stream = spark.readStream.format("ledger").option("endpoint", "https://ledger.internal").load()The hard part is never the code; it is having a position that is durable and ordered. If the upstream system can only tell you "what is new since you last asked", you cannot replay, and no reader can fix that. For the JVM interfaces behind built-in sources see Spark Data Source V2.
Failure modes
- Files overwritten in place. The source tracks paths, so a rewritten file is never re-read, and a replayed batch may read the new bytes. Treat landing files as immutable.
- Partially written files. A writer that streams directly into the watched HDFS directory exposes half-written files to a listing. Write elsewhere and rename.
- Schema drift. With a fixed schema and the default permissive mode, a renamed field becomes NULL and a mistyped one becomes a malformed record. Count nulls on required fields per batch and alert.
- Deleting the checkpoint. The seen-files log goes with it; the restarted query reprocesses every file in the directory.
- Two queries, one directory,
cleanSourceon. One query archives files the other has not read yet. Only enable cleanup when a single query owns the directory. - Huge directories. Listing millions of objects each trigger costs time and request fees. Archive processed files, or partition the landing path by date.
Trade-offs
| Choice | Gain | Cost |
|---|---|---|
| File source | No broker to run, cheap durable input | Listing latency and cost, per-file overhead |
| Kafka source | Low latency, ordered partitions | Retention limits replay, broker operations |
| Table source | Transactional input, time travel | Bound to one table format's rules |
| Per-trigger caps | Bounded batch size and memory | Longer catch-up after outages |
| Custom source | Any upstream with ordered positions | You own replay correctness |
What to do next
- For each production query, name its source and confirm the source is replayable for the longest outage you plan to survive.
- Set an explicit schema on every file source and alert on nulls in required columns.
- Add
maxFilesPerTriggerormaxBytesPerTriggersized from one batch you have measured. - Replace
Trigger.Oncejobs withTrigger.AvailableNow. - Make producers write immutable, complete files; enable
cleanSourceonly where one query owns the directory. - Use
rate-micro-batchin tests that check windowing and watermark logic; keep socket sources out of production. - Read Structured Streaming State next, since replay correctness and state recovery share the same checkpoint.