A practical look at how cloudFiles and Structured Streaming make file ingestion simple, even at millions of files
Some problems in data engineering feel like they should have been solved years ago. Take file ingestion. A cloud storage bucket fills up with new files every few minutes. You need to read them, validate them, and land them somewhere clean. This sounds simple until you try to do it at scale.
At a few files per hour, a simple cron job or directory listing works fine. At a few thousand files per minute, the same approach falls apart. Listing directories becomes slow. Missing files becomes expensive. Schema drift quietly breaks pipelines. And cloud storage API costs start adding up in ways that surprise finance teams.
This is the gap that Auto Loader in Databricks was built to fill. It is not a new idea. It is a better implementation of an old one: let the files tell you when they arrive, instead of asking the storage system over and over again. Today, Auto Loader is a core component of Databricks Lakeflow, which is generally available and provides unified ingestion and orchestration under Unity Catalog.
Most ingestion pipelines start the same way. You point Spark at a directory and read everything in it.
# The approach that works until it does not
df = spark.read.format("parquet").load("/mnt/data/incoming/")This works beautifully for batch jobs on static datasets. But for streaming or near-real-time ingestion, it creates two hidden costs.
First, every trigger has to list the entire directory to find what is new. As the file count grows, this listing time grows with it. The job spends more time figuring out what to read than actually reading.
Second, there is no built-in way to remember what you have already processed. You either re-read everything, which is wasteful, or you build your own bookkeeping, which is fragile. Most teams end up with a hybrid: list files, compare against a watermark, track state in a database or Delta table. It works, but it is plumbing, not business logic.
Auto Loader uses a different approach called cloudFiles. Instead of polling the directory, it asks the cloud provider to notify it when new files arrive. The notification queue becomes the source of truth. The pipeline does not list. It listens.
Here is what the same ingestion looks like with Auto Loader:
from pyspark.sql.functions import col
# Auto Loader reads new files as they arrive
df = spark.readStream \
.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", "/mnt/data/checkpoints/schema") \
.load("/mnt/data/incoming/")
# Write to Delta Lake with exactly-once guarantees
df.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation", "/mnt/data/checkpoints/ingest") \
.table("bronze.raw_events")The key difference is subtle but important. The readStream with cloudFiles does not scan the directory on every trigger. It consumes file arrival notifications from the cloud provider, which means it knows about new files almost immediately and never needs to list the ones it has already seen.
This is the pattern that makes Auto Loader scale from hundreds of files to millions without a change in code.
Auto Loader combines two mechanisms to make this work.
When you start a stream, it first does a one-time backfill using directory listing to catch any files that arrived before the stream began. After that, it switches to notification-based discovery. The cloud provider, whether it is AWS S3, Azure Blob Storage, or Google Cloud Storage, publishes events to a queue whenever an object is created. Auto Loader subscribes to that queue and processes files as notifications arrive.
If notifications are temporarily unavailable, Auto Loader falls back to listing. It is resilient by design. The schema location and checkpoint location handle all the state management, so you do not need to build your own watermark tracking.
# The schema location stores inferred schema over time
# The checkpoint location stores exactly-once stream stateThis dual-mode operation is what makes Auto Loader reliable in production. It is fast when notifications work, and safe when they do not.
One of the hardest parts of long-running ingestion pipelines is handling schema changes. A downstream team adds a column. An upstream system starts sending nested JSON instead of flat records. These changes are normal, but they can break pipelines that assume a fixed schema.
Auto Loader handles this through schema inference and evolution. When you enable it, Auto Loader watches for new fields and either adds them to the inferred schema or stores the raw data in a rescue column, depending on your configuration.
df = spark.readStream \
.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", "/mnt/data/checkpoints/schema") \
.option("cloudFiles.schemaEvolutionMode", "addNewColumns") \
.option("cloudFiles.inferColumnTypes", "true") \
.load("/mnt/data/incoming/")The schemaEvolutionMode option controls what happens when a new field appears. The addNewColumns mode merges the new field into the schema automatically. The rescue mode captures unexpected fields in a JSON blob so you do not lose data even when the schema diverges.
This is not magic. It is a disciplined approach to a common problem. Schema evolution still needs monitoring and governance. But Auto Loader gives you the machinery to handle it without writing custom parsers or restarting your stream.
There is a practical reason Auto Loader matters beyond convenience. Cloud storage APIs charge per request. Every LIST operation counts. Every HEAD request counts. At scale, a polling-based ingestion job can generate millions of API calls per day.
Auto Loader reduces this dramatically because it stops polling. Once the notification stream is active, the only API calls are the actual file reads. For pipelines processing millions of files, this difference can be measured in real dollars.
The exact savings depend on your cloud provider, your file arrival rate, and your current polling frequency. But the direction is consistent: fewer API calls, lower costs, and a simpler architecture.
Auto Loader is not the right tool for every ingestion problem. It shines in specific situations.
It works well when new files arrive continuously and you want to process them incrementally. It works well when you need exactly-once semantics and automatic checkpointing. It works well when you want schema inference and evolution without building your own metadata store.
It is less suitable for one-time bulk loads of static datasets, where a simple batch read is faster and simpler. It is also not ideal when you need complex file-level filtering or custom parsing logic that cannot be expressed in standard Spark readers. For more complex source systems, Lakeflow Connect now provides GA query-based connectors for over 100 sources.
Understanding these boundaries is part of thinking like a data engineer. The goal is not to use every feature. It is to match the right tool to the right pattern.
Here is a practical Auto Loader setup that handles JSON files landing in cloud storage, with schema evolution and a Delta Lake sink:
from pyspark.sql.functions import current_timestamp, input_file_name
# Read new JSON files as they arrive
raw_stream = spark.readStream \
.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", "/Volumes/bronze/checkpoints/schema") \
.option("cloudFiles.schemaEvolutionMode", "addNewColumns") \
.option("cloudFiles.inferColumnTypes", "true") \
.option("cloudFiles.maxFilesPerTrigger", "1000") \
.load("s3a://landing-bucket/events/")
# Add standard audit columns
enriched = raw_stream \
.withColumn("_ingested_at", current_timestamp()) \
.withColumn("_source_file", input_file_name())
# Write to bronze layer with exactly-once guarantees
enriched.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation", "/Volumes/bronze/checkpoints/ingest") \
.option("mergeSchema", "true") \
.trigger(availableNow=True) \
.table("bronze.raw_events")The availableNow=True trigger processes all currently available files and then stops. For continuous ingestion, you would use a time-based trigger instead.
# For continuous ingestion, replace the trigger
.trigger(processingTime="1 minute")This gives you a pipeline that scales with your file volume, handles schema drift, and lands data in Delta Lake with transactional guarantees. The code is small because the platform handles the hard parts.
Auto Loader is a useful tool, but the larger lesson is about architecture. Modern data platforms are not just about storage and compute. They are about the automation layer that sits between your data sources and your analytics.
With the GA release of Lakeflow, these patterns are now part of a unified ingestion, transformation, and orchestration framework. A well-designed ingestion pipeline should:
Auto Loader addresses all four of these concerns for file-based ingestion. It is a pattern worth understanding even if you are not using Databricks today, because the ideas, notification-based discovery, schema evolution, and checkpointed exactly-once processing, are becoming standard across platforms.
Auto Loader connects naturally to several topics covered in BricksNotes. The Data Sources lesson explains how Spark reads different file formats. The Streaming lesson covers Structured Streaming in detail, including triggers, checkpoints, and fault tolerance. The Delta Lake lesson shows how the transaction log enables reliable incremental writes. And the Partitioning & Performance lesson helps you think about file layout once your data lands.
If you are building ingestion pipelines that need to grow, start by asking what the platform can handle for you. Often, the simplest code is the most scalable.