11 min read

Auto Loader Does Not Scale by Default

Auto Loader's default detection mode is a full directory listing, and the checkpoint never changed that. A practical look at the five defaults that decide whether your ingestion scales, and the file count that tells you which mode you need.


I've lost count of how many Auto Loader streams I've seen that look like this:

(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", schema_path)
  .load(landing_path))

Four lines, works on day one, and everyone moves on. I've written that exact block myself more than once. The problem is that nothing in it tells you how it's going to behave at ten thousand files, or at ten million. Those are two different products wearing the same syntax, and the difference comes down to defaults you never typed.

This is the version I wish someone had drawn for me on a whiteboard before I put Auto Loader on a telemetry feed that produced a few hundred thousand small files a month.

Auto Loader has two ways of finding files, and the default is the slow one

When the stream ticks, Auto Loader has to answer one question: what's new since last time? It can do that in two ways.

Directory listing. This is what you get out of the box. Every trigger, it lists the whole input path, then compares what it found against what it already knows about. If your landing zone has 600k objects in it, every trigger enumerates 600k objects, finds the 2k that are new, and throws the rest away. The cost of that listing has nothing to do with how much data you're ingesting. It tracks how many objects are sitting in the path, which, unless you archive, is every object that ever landed there.

File notification. Cloud storage emits an event per new object into a queue, and Auto Loader drains the queue. It never lists the directory on a normal run. Cost tracks new files only, flat regardless of how big the path has grown. On Azure this used to mean wiring up Event Grid and a Storage Queue yourself; the current path is file events on the Unity Catalog external location, which Databricks sets up for you.

Both modes read each file exactly once. That part is identical. The only thing that differs is what it costs to discover the file in the first place, and that's precisely the part people assume Auto Loader has already solved for them.

01 two detection modes

One combination I'd call out specifically, because I've seen it in production: a continuously running trigger in directory listing mode. That is the worst case on both axes. You're paying for a full listing over and over to find a handful of files. If you don't actually need sub-minute freshness (most of us don't), a scheduled job with Trigger.AvailableNow is still a stream, still uses the checkpoint, and costs a fraction of it.

The checkpoint is a ledger, not an index

This is the misunderstanding that costs the most money, so I'll spend a paragraph on it.

Inside the checkpoint location there's a RocksDB store. As Auto Loader discovers files it writes their metadata there, and on every run it checks that store before processing anything. If a file is in it, it's skipped. That's where the "each file once" guarantee comes from, and it's real.

But look at where in the sequence it sits. In directory listing mode the order is: list everything, then consult the ledger, then discard what's already there. The listing bill is paid before the ledger is ever opened. The checkpoint stops the re-read. It does nothing whatsoever about the re-listing.

02 checkpoint is a ledger

So the typical story goes: a team has a spark.read over a growing folder, it gets slow, someone says "use Auto Loader, it's incremental". They switch, keep the default mode, and the job is exactly as slow as before, because the slow part was the listing and the listing didn't change. The fix they needed was file notification. The checkpoint was never the part that was going to help.

While we're on guarantees, one boundary. Exactly-once covers committed micro-batches written straight to a Delta table. It does not cover foreachBatch. That path is documented as at-least-once, and the only thing that makes a replayed batch a no-op is setting txnAppId and txnVersion (bound to the batch id) on the write. And if you ever start a fresh checkpoint, give it a fresh txnAppId too, because batch ids restart at zero and Delta will quietly skip everything you just replayed thinking it has already seen it.

File count picks the mode. Not gigabytes.

When someone tells me "it's about 2 TB a day" I still don't know which mode they need. Two terabytes as a couple of thousand big Parquet files is nothing to list. The same two terabytes as a few million small JSON objects, which is what telemetry and API dumps tend to look like, is a listing problem on day one.

The Databricks line is roughly: thousands of files over time, COPY INTO or directory listing is fine; millions, you want file notification. With file events it's built to handle on the order of millions of files per hour. Read your numbers as arrival rates, though. Two thousand files a day takes years to become a problem. Three million a day is one before lunch.

Two things to know before you flip the switch:

There's no middle option any more. Incremental listing, the half-measure some of us still remember, is deprecated. It's listing or notification.

The 7-day rule. On file events, the stream has to run at least once every 7 days. Miss that and the next run falls back to a full directory listing. A job on a weekly schedule is sitting right on that line; the first time it slips (a failed run, a paused job during an incident, a holiday) it silently pays the exact bill you switched modes to avoid.

And a note on the most misleadingly named option in the whole product. cloudFiles.backfillInterval sounds like a rewind. It isn't. It schedules an asynchronous full listing so Auto Loader can pick up files the event queue might have dropped. It points forward at discovery, never backward at reprocessing, and it will never re-emit a file already in the ledger. On file events it isn't even a supported option, because Databricks runs that reconciliation for you.

Schema drift: only one thing ever stops the stream

Someone upstream adds a column on a Monday morning and your stream goes red. People read that as a failure. It's actually the handoff.

If you didn't pass a schema, Auto Loader infers one into the schema location and cloudFiles.schemaEvolutionMode defaults to addNewColumns. On a new top-level column it writes the updated schema to that location and then fails the query with UnknownFieldException. The next run starts from the new schema and the column is just there.

The catch is that Auto Loader doesn't restart itself. The intended setup is a Lakeflow job with automatic restart on. With that, a new column costs you one failed run. Without it, the "evolution" you were promised is a red stream waiting for someone to notice.

And there's a second catch that bit me. Auto Loader evolves the schema it reads with. Your Delta sink still enforces the schema it writes to. So the restart lands, the read succeeds, and the write rejects the new column. Unless mergeSchema is on the write (or auto-merge is on for the session), you get the same job failing on the write side instead of the read side, with the retries spacing out a little more each time.

The other modes, so you choose instead of inherit:

  • failOnNewColumns: fails and stays down. A human has to update the schema or remove the file.
  • rescue: never evolves, never fails. New columns land in _rescued_data.
  • none: never evolves, never fails, and ignores the column entirely. This is the default the moment you supply your own schema, and under none nothing is rescued unless you explicitly set rescuedDataColumn.

03 schema drift tree

That last point matters more than it looks. Two defaults, opposite postures, same product. Infer the schema and you get a loud failure on a new column plus a rescue column for free. Supply a schema and you get silence, dropped columns, and no _rescued_data at all. I've seen a drift alert built on _rescued_data against a provided-schema stream. It was querying a column the pipeline never materialised. It was never going to fire.

What about type changes? On the four modes above, nothing fails. A value that won't parse into the tracked type gets rescued: the typed column goes null, the raw value lands in _rescued_data (if you have one). The only exception is addNewColumnsWithTypeWidening, which is opt-in and still in preview; it treats a widenable type change like a new column, widens the schema, and fails the query so the value ends up in the typed column.

For JSON, CSV and XML there's a prior problem: cloudFiles.inferColumnTypes is off by default, so every column is inferred as a string, nested fields included. If everything's a string there's no type to mismatch, so nothing ever gets rescued for a bad value. The quality rule downstream that compares reading_kwh to a number isn't doing what it looks like it's doing. Turn the option on, or supply a typed schema.

"Yesterday's files were bad, reprocess just that day"

You can't, not through Auto Loader. There is no per-day rewind. Once a file is in the ledger, Auto Loader is done with it. You have three options and they're nowhere near equivalent.

  1. cloudFiles.allowOverwrites. Reprocesses a file when it's overwritten or appended to. The docs are explicit that the resulting duplicates are your problem.
  2. A fresh checkpoint location. Abandons the old stream and starts a new one. With cloudFiles.includeExistingFiles at its default of true, that means every file in the path, not yesterday's. On a 600k-object landing zone that's a very expensive way to redo one day.
  3. Scope it yourself. A one-off batch read of that day's partition, and an idempotent write: replaceWhere over that day, or a MERGE on the key if the corrected rows carry the same keys. The replacement replaces instead of doubling.

The third is almost always the right answer and almost nobody reaches for it, because it feels like going around the streaming abstraction. It isn't. The abstraction was built for forward motion. Fixing the past is a batch job, and that's fine.

The five defaults you inherit with the snippet

Copying a working snippet means inheriting all of these. Each one is probably something you want to override, and each has a consequence if you don't.

SettingDefaultWhat it does to youWhat I'd set
Detection modedirectory listingLists the whole path every triggerFile notification via file events, and make sure it runs well inside 7 days
cloudFiles.schemaEvolutionModeaddNewColumns (inferred) / none (supplied)Loud failure vs. silent dropped columns, from the same productPick one on purpose; if supplying a schema, set rescuedDataColumn
cloudFiles.inferColumnTypesfalseJSON/CSV/XML all strings, so nothing is ever rescued for a bad valuetrue, or a typed schema
cloudFiles.maxFilesPerTrigger1000Backlog drains 1000 files per batch after an outageRaise it deliberately at high file counts; the unprefixed maxFilesPerTrigger is silently ignored on a cloudFiles read
cloudFiles.includeExistingFilestrueA new checkpoint means "everything", not "from now"Know it before you ever touch a checkpoint

One more that isn't a default but belongs in the same breath: cloudFiles.maxBytesPerTrigger isn't set at all. The throttle you have is a file-count throttle, nothing else.

The sentence for the design review

Discovery and processing are separate problems. The ledger solves processing, in either mode, so no file is read twice. Listing cost is a discovery problem with two real levers: stop listing (file notification), or stop the path from growing (archive processed files out of it). Triggering less often pays the bill less often but doesn't shrink it. File count picks the mode, never data volume. And apart from the type-widening preview, schema drift has exactly one failure trigger, a new column; everything else lands silently in a rescued column you may not even have.

There are two answers to "how does your ingestion scale?" One is "we use Auto Loader, it's incremental". That's the snippet talking, and it survives until someone asks what the file count is. The other is something like "we're on directory listing at around six hundred thousand objects growing forty thousand a day, fine for now, we move to file notification before we hit a few million, evolution mode is addNewColumns under a job that restarts, mergeSchema on the sink, so a new column costs us one failed run."

The distance between those two answers isn't years of experience. It's five defaults and knowing your own file count.

So, what's the file count on your busiest landing path, and which mode is it on? I'd bet most are still on the default. Better to find out now than in a review.

Comments

No comments yet. Be the first to leave one below.

Leave a comment

Comments are reviewed before appearing.