Auto Loader is a Structured Streaming source called cloudFiles that picks up new files from cloud storage as they arrive. It remembers which files it has processed, so each file is ingested once, even with millions of files in the folder.
Basic use
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "/chk/orders/schema")
.load("s3://company-landing/orders/"))
(df.writeStream
.option("checkpointLocation", "/chk/orders")
.trigger(availableNow=True)
.toTable("bronze.orders"))How it avoids reprocessing
It records the files it has discovered and processed in the checkpoint, in a RocksDB key-value store. On the next run it processes only files not yet recorded. You do not have to track file names yourself or move files after processing.
Two ways to find new files
- Directory listing mode (the default): Auto Loader lists the folder to find new files. It is easy to set up, and works well for modest numbers of files, with incremental listing that avoids re-listing everything.
- File notification mode: it sets up cloud event notifications and a queue (for example S3 events with SQS), and reads new file names from there. This scales much better for very large folders, at the cost of setting up permissions and cloud resources.
Schema inference and evolution
With schemaLocation, Auto Loader infers the schema from a sample and stores it. When new columns appear in later files, by default the stream stops and picks up the new schema on restart (evolution mode addNewColumns). Data that does not fit the schema goes to a _rescued_data column rather than being lost, so you can inspect it.
Triggers
For continuous ingestion, leave it running. For cheaper, scheduled ingestion, use the availableNow trigger: it processes everything new and then stops, so a job can run hourly and the cluster only exists while there is work.
What to say
Compared with a plain spark.read on a folder, Auto Loader is incremental, handles huge file counts, and deals with schema drift. Compared with Snowpipe, it is the Databricks equivalent for file-based ingestion.