Incremental file ingestion
Files keep landing in a folder: exports, logs, partner drops. The job is to load each file exactly once, keep up as volume grows, and survive schema changes. That is incremental ingestion.
On this page
Read first: Structured Streaming basics · Data lake storage: formats and layout (skip if you know them)
The problem
The naive job lists the landing folder, compares with what it loaded last time (by file name or modification time), and loads the difference. It works on day one. Then the folder holds millions of files and listing takes longer than loading; a file is written slowly and picked up half-written; a job fails between loading and recording; two runs overlap. Each failure means duplicated or missing data.
The file source as a stream
Structured Streaming treats a folder as a stream of files. The checkpointcheckpoint: In streaming, the folder where Spark records which input it has processed and its state, so a restarted query continues where it stopped. Learn more → records which files each micro-batch processed, so each file is processed once, even across restarts. Run it continuously, or as a scheduled batch job with an availableNow trigger, which processes everything new and stops:
(spark.readStream .format("json") .schema(order_schema) # streams need a schema up front .option("maxFilesPerTrigger", 1000) # bound each micro-batch .load("s3://landing/orders/") .writeStream .format("delta") .option("checkpointLocation", "s3://chk/orders_bronze/") .trigger(availableNow=True) # Spark 3.3+: process what is new, then stop .toTable("bronze.orders"))
- Exactly once comes from the checkpoint plus a transactional sink such as Delta or Iceberg.
- availableNow gives batch scheduling with streaming bookkeeping: no "what did I load last time" table to maintain. It also respects
maxFilesPerTrigger, splitting a large backlog into several batches. - cleanSource (
archiveordelete) moves or deletes processed files, which keeps the landing folder small. - The catch: the source still lists the directory every batch. With millions of files that listing becomes the bottleneck.
Auto Loader and COPY INTO
On Databricks, two tools address the gaps. They are proprietary features of that platform, but the ideas show up in interviews generally.
Auto Loader (cloudFiles) | COPY INTO | |
|---|---|---|
| Interface | A streaming source: spark.readStream.format("cloudFiles") | A SQL command run as a batch |
| Discovering files | Directory listing, or file notification mode using cloud events and a queue, which avoids listing entirely | Lists the source path each run |
| Tracking | Checkpoint with a scalable key-value store of processed files | Table metadata records loaded files; reruns skip them |
| Schema | Infers and stores the schema, evolves it when new columns appear, and puts unexpected values in _rescued_data | Explicit or inferred per run |
| Fits | Continuous or frequent ingestion, millions of files | Simple, occasional loads of thousands of files |
Schema drift and bad records
New columns
Type changes and junk
Raw first
Backfills and late files
A stream only processes files it has not seen. To reprocess history (a parsing bug, a new column), start a new checkpoint into a new or truncated target, or run a separate batch backfill over a date range and keep the stream for new data. Never delete the checkpoint of a running pipeline casually: it will reload everything into the same target and duplicate it.
ingest_date=2025-03-01/) even if the data has its own event dates. Late-arriving files then simply land in today's folder, and backfills can target specific ingestion dates.Common mistakes
- Tracking loaded files by modification time — Clock skewskew: When one key has far more rows than the others, so one task does most of the work while the rest wait. Learn more →, slow writers and copied files break it. Use a source that records processed files.
- Picking up files while they are still being written — Writers should write to a temporary name or folder and rename or move when done.
- Deleting the checkpoint to "fix" a stream — Everything is reprocessed into the same target: duplicates.
- Inferring the schema on every batch run — The schema changes with whatever files happen to arrive.
Interview prep
THE QUESTION
"Partners drop CSV files into S3 all day. Design an ingestion that loads each file exactly once into a Delta table."
Avoid saying: "store the last run time and load files modified after it". Clock skew, slow writers and copied files make it skip or double-load files.
What the interviewer asks next. Answer out loud first, then open the strong answer.
Follow-up"What is the difference between trigger(once=True) and trigger(availableNow=True)?"
maxFilesPerTrigger, which can be huge after a backlog. availableNow (Spark 3.3+) splits the backlog into several bounded batches, and replaces once, which is deprecated.Scenario"A parsing bug corrupted the last two weeks in bronze. How do you reload those files?"
Trap"The file source detects when an existing file is overwritten and reloads it."
What you learned
- Why "list the folder and load what is new" breaks down
- How the Structured Streaming file source tracks processed files
- What Auto Loader and COPY INTO add, and where they apply
- How to handle schema drift, bad records and backfills
Key takeaways
- A streaming file source with a checkpoint loads each file exactly once.
- trigger(availableNow=True) gives scheduled batch runs with streaming bookkeeping.
- Auto Loader adds notification-based discovery and schema evolution; COPY INTO is an idempotent batch load (both Databricks).
- Land raw first, keep bad records visible, and backfill with a new checkpoint.
Check yourself
3 questionsWhat makes the file source process each file only once across restarts?
Show the answer
The checkpoint log of processed files, plus a transactional sink. The checkpoint records processed files per batch; the sink commits atomically.
What does trigger(availableNow=True) do?
Show the answer
Processes all new data, possibly in several batches, then stops. It turns a stream into an incremental batch job.
Why can a plain file stream slow down as a landing folder grows?
Show the answer
Each batch lists the directory, and listing millions of files is slow. Listing cost grows with the folder; archive processed files or use notifications.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 18 of 30 · 4 min readPartitioning done right
When to partition a table, how to choose the column, and how over-partitioning backfires.
Related lessons
Structured Streaming basicsData lake & lakehouse · 5 min read
Streaming into Delta tablesData lake & lakehouse · 4 min read
The medallion architecture
Previous: Streaming into Delta tables
Primary sources: Structured Streaming: file source · Databricks Auto Loader