- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
04-01-2026 05:46 AM
Adding on to @edonaire, which are accurate.
@beaglerot , your contacts project is the right use case for the pattern you have. Small data, infrequent changes, direct read into bronze. That works. The real question you're asking is what happens when the data gets bigger and changes faster. Here's how I'd think about it.
There are two viable patterns, and the right one depends on what you need from your raw layer.
Option A: Direct to bronze via the Data Source API (no JSON landing zone)
If you're comfortable treating your bronze Delta table as the system of record, skip the intermediate JSON files entirely. Include the full API payload (or a raw_payload column) in your bronze table so you still have an immutable representation of what the API returned. Recovery and reprocessing come from Delta time travel and cloning rather than re-reading raw files.
@dp.table(
name="contacts_bronze",
comment="Raw contacts from Google People API"
)
def contacts_bronze():
return (
spark.readStream
.format("google_people")
.option("scope", "google-api")
.load()
.withColumn("ingested_at", current_timestamp())
)
This is the simplest mental model and fewest moving parts. It works well when your bottleneck is API throughput rather than storage, and when you're wiring into DLT or Workflows.
One note: this requires your custom data source to implement the SimpleDataSourceStreamReader class. If your source only supports batch reads, you'd use spark.read with a scheduled job instead of readStream.
Option B: Land raw JSON to a volume, then Auto Loader into bronze
If you need a file-level audit trail, expect upstream schema changes, or have multiple downstream teams consuming the same raw feed in different ways, land the data as files first.
df = (
spark.read.format("google_people")
.option("scope", "google-api")
.load()
.withColumn("ingested_at", current_timestamp())
.withColumn("ingest_date", to_date("ingested_at"))
)
(df.write
.mode("append")
.partitionBy("ingest_date")
.format("json")
.save("/Volumes/raw/google_people"))
Then point Auto Loader at that path:
bronze = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/Volumes/raw/google_people")
)
A few things on file structure if you go this route:
Avoid one tiny file per pull. Aim for fewer, larger files (tens to hundreds of MB per batch) to keep Auto Loader and downstream queries efficient. Partition by date using a path like /Volumes/raw/google_people/ingest_date=YYYY-MM-DD/ rather than encoding timestamps only in filenames. Add ingestion metadata as columns: ingested_at, source_system, api_version, and optionally a batch_id for traceability.
How to decide between them?
Use Option A when you want the simplest architecture and you're comfortable with Delta as your recovery mechanism. The question to ask yourself: "If something breaks downstream, can I reprocess from Delta time travel, or do I need the original API payloads sitting in storage?"
Use Option B when the answer to that question is "I need the raw payloads," or when you have strict audit requirements, or when you don't want to hit the API again to reprocess.
In both cases, the main win of the Python Data Source API is the same: standardizing the edge. One well-tested connector that handles auth, pagination, retries, and schema, instead of N slightly different notebooks all doing their own version of that logic. It doesn't replace your ingestion architecture. It standardizes the extraction layer that feeds into it.
Cheers, Lou