Skip to content
Good engineers know 6 min read · All formats 2 practice problems ↓

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.

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:

PySpark · Incremental load of JSON drops into a Delta table
(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 (archive or delete) 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
InterfaceA streaming source: spark.readStream.format("cloudFiles")A SQL command run as a batch
Discovering filesDirectory listing, or file notification mode using cloud events and a queue, which avoids listing entirelyLists the source path each run
TrackingCheckpoint with a scalable key-value store of processed filesTable metadata records loaded files; reruns skip them
SchemaInfers and stores the schema, evolves it when new columns appear, and puts unexpected values in _rescued_dataExplicit or inferred per run
FitsContinuous or frequent ingestion, millions of filesSimple, occasional loads of thousands of files

Schema drift and bad records

New columns

With an explicit schema, new fields are ignored, not lost forever if bronze keeps the raw payload. With Auto Loader's schema evolution, the stream stops on a new column, updates the stored schema, and continues after a restart.

Type changes and junk

Values that do not fit the schema become null in permissive mode. Keep them visible: a corrupt-record or rescued-data column, counted and alerted on, rather than silently null.

Raw first

Land data in bronze with minimal parsing (or as a string or VARIANT payload), and parse strictly in silver. A schema problem then never blocks ingestion.

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.

Tip: land files by ingestion date folders (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."

Use the Structured Streaming file source (or Auto Loader on Databricks) with an explicit schema, a checkpoint, and a Delta sink; schedule it with trigger(availableNow=True) if near-real-time is not needed. The checkpoint tracks processed files and Delta commits atomically, giving exactly-once. Land into bronze with ingestion metadata (file name via _metadata, load time), keep bad records visible, and parse into silver. Mention listing cost at scale (file notification mode, cleanSource archiving), files written in place, and how to backfill with a new checkpoint.

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)?"
Both process what is new and stop. Once processes everything in one micro-batch, ignoring 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?"
The stream will not reprocess files its checkpoint has seen. Run a batch backfill that reads those ingestion-date folders, and replace the affected rows in the target (overwrite by ingestion date, or delete and re-insert), while the stream keeps loading new files. Or start a new stream with a new checkpoint into a fresh table and swap.
Trap"The file source detects when an existing file is overwritten and reloads it."
By default it tracks files by path. A file rewritten in place under the same name is not processed again. Producers should write new files with unique names.

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 questions

What 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.

Solve: Late Data in the Raw Zone →

Keep going

Up next · lesson 18 of 30 · 4 min read
Partitioning done right
When to partition a table, how to choose the column, and how over-partitioning backfires.

Related lessons

Previous: Streaming into Delta tables

Primary sources: Structured Streaming: file source · Databricks Auto Loader