Technical deep dive
From Airflow DAGs to Metadata-Driven Aviation Pipelines
by Ondrej Grünwald, Founder / CEO
How SkyAlgorithm replaced hand-written Airflow DAGs with metadata-driven aviation pipelines for lineage, versioning, previews, schema drift handling, and safer reruns.
The first version of the SkyAlgorithm data platform was a familiar shape: raw aviation data landed in the lake, Airflow DAGs picked up hourly partitions, DuckDB loaded them into Iceberg, and curators promoted silver rows into gold tables.
That worked. It also left too much of the platform's meaning in places that were hard to inspect.
One raw ADS-B field made the problem visible. typecode was present in the raw files, but it never reached silver because the loader's closed projection did not name it. The job still ran green. The table still looked healthy. Nothing in Airflow could say "the source shape changed and the transform ignored it," because Airflow only knew that a task succeeded.
A DAG file can tell you what command runs. It does not naturally tell you which logical dataset the command reads, which version of the transformation produced a past run, whether a schema change was expected, whether a replacement write can rebuild everything it deletes, or which monitoring checks should move when the job is replaced. Those questions matter once a platform has real history behind it: old jobs, replaced jobs, backfills, schema changes, monitoring checks, and operational data products people rely on.
This post is about the layer we built after the lakehouse: a metadata-driven pipeline model for aviation data jobs. Airflow still schedules and executes work, but the pipeline itself is no longer a hand-written DAG. It is a stored definition with versions, dataset references, run records, lineage, preview, and explicit schema behavior.

What the migration had to prove
Same output
A definition reading raw ADS-B had to reproduce the established silver rows for the same hour before it could write the real target.
Bounded writes
A replacement run had to delete only the hour it could rebuild, then write that hour back without duplicating or emptying the extent.
Moved ownership
Catalog ownership and monitoring checks had to move with the pipeline, or the old DAG would still be treated as the source of truth.
1. The Problem with "Just Add a DAG"
The lakehouse post covered the main platform shape: raw ADS-B, FlightAware, METAR, and TAF data land in MinIO; silver and gold Iceberg tables are queried through Trino; TimescaleDB serves the hot path; and Airflow runs the batch jobs.
The first Airflow DAGs were direct and useful:
- load one hour of raw ADS-B into silver
- load one hour of FlightAware snapshots into silver
- curate ADS-B silver into gold
- curate FlightAware silver into gold
- run maintenance over Iceberg tables
Each DAG wrapped a script with environment variables. For a small number of jobs, this is hard to beat. The DAG is visible, deployable, and familiar.
The trouble starts when the platform needs to answer operational questions about the data rather than about the scheduler.
Which dataset did this job mean to read? Which physical prefix did it resolve at runtime? Which version of the transform wrote this hour? Did it delete before it inserted? Did a new source field appear? Did the table widen during the run? If this DAG is paused, which freshness checks now need to point at the replacement?
Airflow has answers for task state. It does not know the domain contract of an aviation data product unless we encode that contract somewhere else.
That "somewhere else" became the pipeline definition.
2. A Pipeline Is a Stored Definition
A pipeline definition is stored data. It lives in Postgres as JSONB, has immutable versions, and is executed by a generic runner.
At the simplest level, the shape is:
Source -> Transform -> Target
The source says what data the run reads. The transform says what rows should become. The target says where rows land and how they are written. The schedule says when a window is due.
For example, an hourly ADS-B silver load is not "the DAG that runs this Python script." It is a definition that says:
{
"name": "adsb_silver_hourly",
"kind": "batch",
"sources": [
{
"name": "raw",
"dataset": "adsb.raw",
"window": {
"kind": "partition",
"granularity": "hour",
"lag": "PT10M"
},
"options": {
"schema_mode": "full",
"on_parse_error": "skip"
},
"schema_policy": {
"on_new_field": "warn",
"on_missing_field": "null"
}
}
],
"transform": {
"dialect": "duckdb",
"sql": "SELECT ... FROM raw"
},
"target": {
"dataset": "adsb.silver",
"write_mode": "append",
"on_schema_change": "fail"
},
"schedule": {
"cron": "10 * * * *",
"catchup": false
}
}
The ellipsis in the SQL is the boring part. The important part is everything around it.
The definition names a logical source and target. It says the run reads one hourly partition with a 10-minute lag. It declares whether source schema drift is a warning or a failure. It declares whether the target is allowed to evolve. It has a schedule, but saving it does not imply it should start running.
Those details used to live across DAG files, scripts, config, table DDL, and memory. Once they are in the definition, the platform can reason about them.
There is a boundary here that matters: the definition stores the parts of the job the platform must validate or explain later. It does not store everything Airflow, Iceberg, or MinIO already own.
The pipeline model owns dataset identity, source windows, schema policy, write mode, target behavior, versions, and run records. The catalog owns the target binding and lifecycle. Iceberg owns table schema, partitioning, files, and snapshots. Airflow owns scheduler mechanics such as retries and executor behavior. MinIO owns object listings.
That separation keeps the metadata useful instead of turning it into a stale copy of every system underneath it.
Replacement writes need an extent
The riskiest write mode is not append. It is replace_window, because it deletes existing rows before writing a bounded slice back.
That is why a replacement target has to say which target rows a run owns:
{
"target": {
"dataset": "adsb.silver",
"write_mode": "replace_window",
"replace": {
"column": "ingest_ts",
"granularity": "hour",
"column_type": "epoch_s"
},
"on_schema_change": "fail"
}
}
This is deliberately more specific than "replace the partition." Raw ADS-B is read by ingest-hour partitions, while the cleaned table also carries event-time columns. Deleting by the wrong clock can remove rows the run did not read and therefore cannot rebuild.
The invariant is simple: a run must be able to rebuild, from its declared source window, every target row it deletes. If the definition cannot state that relationship, the job should not be a replacement write.
3. Logical Datasets, Not Bucket Paths
The first design rule was simple:
Pipeline sources reference datasets, never locations.
The source in the definition is adsb.raw, not a concrete hourly prefix such as lake/topics/adsb.raw/year=2026/month=08/day=24/hour=10/. The target is adsb.silver, not the physical table name iceberg.aviation_silver.adsb_clean.
The difference is identity versus location. adsb.raw means "the raw ADS-B dataset" regardless of where it currently lives. The hourly prefix is one physical slice chosen for one run. adsb.silver means "the cleaned ADS-B dataset" regardless of the Iceberg schema and table that currently implement it.
That is what makes the model useful. A pipeline definition can stay stable while the catalog binding changes, and a run can still record the exact bucket prefix or table it touched.
Raw ADS-B currently binds to an object-store prefix:
{
"id": "adsb.raw",
"layer": "raw",
"binding": {
"type": "object-store",
"bucket": "lake",
"prefix": "topics/adsb.raw/",
"format": "json",
"compression": "gzip",
"partitioning": "hive",
"partition_keys": ["year", "month", "day", "hour"]
}
}
Silver ADS-B binds to an Iceberg table:
{
"id": "adsb.silver",
"layer": "silver",
"binding": {
"type": "iceberg",
"catalog": "iceberg",
"schema": "aviation_silver",
"table": "adsb_clean"
}
}
The definition says what it meant. The run records what that meant physically at execution time: which bucket, which prefix, which table, and which resolved window.
That split matters because storage moves, defaults drift, and deployment profiles differ. One old loader path carried a default bucket name that had never existed in production and only worked because the runtime environment overrode it. A pipeline definition keyed on a bucket path would have turned that kind of mistake into identity. A logical dataset keeps the contract stable while the binding remains editable.
It also makes lineage cheap. If a definition reads adsb.raw and writes adsb.silver, the lineage edge is already there. There is no separate system that has to parse logs and infer what moved.
There is another operational effect: a pipeline can only write datasets the catalog marks as managed by this layer. During the cutover, targets moved from external to managed only after the replacement definitions had run cleanly. That made "who writes this dataset now?" a catalog fact rather than a convention hidden in Airflow.
4. Versions Make Re-Runs Explainable
Editing a pipeline never rewrites an existing version. It inserts a new one.
That is a small database rule with a large operational effect. A run points at the exact version it executed, so a later edit does not change the meaning of yesterday's run history.
Without versioning, every historical run becomes approximate. You can see that adsb_silver_hourly ran at 10:10, but if the definition has changed since then, what exactly ran? Which column list? Which write mode? Which source policy? Which target policy?
With immutable versions, the answer is mechanical:
pipeline_runs.run_id
-> pipeline_definition_version.id
-> stored definition JSON
That does not make every pipeline correct. It makes incorrect output debuggable.
For aviation data, this matters because the same hour may be reprocessed several times. ADS-B data can arrive late. Raw partitions can be replayed from an edge receiver. A transform may need to be widened after a new field appears. The system has to let an operator rerun a bounded window and still explain why the output changed.
5. Preview Before Writing
The most dangerous pipeline is not the one that fails validation. It is the one that validates, runs on schedule, and writes the wrong window.
That is why authoring moved through a builder instead of a free-form JSON editor. In practice, new definitions should be developed against local or staging data first. Production should mostly execute promoted definitions, trigger bounded reruns, and expose enough detail to explain what happened.
The console builder follows the same flow an engineer uses manually:
- Scan a real source window.
- Show the observed fields, types, and non-null counts.
- Preview the transform with the real planner and executor.
- Save the definition disabled.
- Run it once deliberately.
- Add a schedule only after the output is right.
The preview is not a mock. It resolves the dataset binding, lists the raw objects, applies the window, creates the DuckDB view, and runs the transform without writing the target.
That one constraint catches a class of errors that static validation cannot: wrong source binding, wrong partition grain, wrong column name, unexpected raw schema, and SQL that only fails when pointed at real files.
This mattered in the first real pipeline runs. The scratch ADS-B silver definition read one UTC hour, touched 24 raw objects, wrote 480,000 rows, and diffed cleanly against the existing silver output for the same hour. The only expected difference was typecode: it was present in raw and reported as a new source field, but it did not appear in the projection until a new definition version added it deliberately.
That is the behavior we wanted. Schema drift was visible without silently widening the product table, and promotion became a one-field target change after the scratch output proved equivalent.
The console is intentionally conservative here. Enable and disable are cheap reversible controls. Free-form editing of stored production definitions is not the main path because a valid definition can still name the wrong replace window. In local development that may only dirty a scratch table. In production or shared staging, the same mistake can remove the rows a downstream product expects until the window is rebuilt.
6. Schema Drift Is a Runtime Fact
Raw aviation feeds move. New ADS-B fields appear. Providers alter optional fields. A receiver sends a field in one hour and not the next.
The old loader had a fixed projection. If raw carried a new field that silver did not, the pipeline simply did not mention it. That is how typecode could exist in raw data without reaching silver and without an obvious alert.
The definition model separates source schema policy from target schema policy.
For a raw source:
{
"schema_policy": {
"on_new_field": "warn",
"on_missing_field": "null"
}
}
For a governed target:
{
"on_schema_change": "fail"
}
Those are different decisions. A new field in raw may be useful and expected enough to record as a warning. A transform producing a column that the target table does not have is usually a defect unless the definition explicitly allows add_columns.
When on_schema_change is fail, the runner checks the target shape before deleting or inserting. If the transform produces squawk but the target has no squawk column, the run fails before it touches the table.
That ordering is load-bearing. In a replacement run, discovering the mismatch after the delete would leave the target extent empty until somebody noticed and rebuilt it. A schema refusal has to happen while the run is still a read.
When on_schema_change is add_columns, the widening is recorded on the run:
{
"schema_changes": {
"iceberg.aviation_staging.mixed_probe_silver": ["category"]
}
}
That record matters. Otherwise a schema change is only visible afterward as "the table has a new column now." The run history should say which run changed it.
7. Steps Let the Pipeline Match the Data Product
Not every data product is one SQL statement. ADS-B gold now includes multiple pieces: curated positions, emergency squawk events, reconstructed flight tables, geometries, and airport movement candidates.
The definition model supports steps:
{
"steps": [
{
"name": "curate_gold",
"engine": "container",
"inputs": [{ "dataset": "adsb.silver" }],
"outputs": [{ "mode": "materialized", "dataset": "adsb.gold" }]
},
{
"name": "emergency_events",
"engine": "container",
"inputs": [{ "dataset": "adsb.gold" }],
"outputs": [{ "mode": "materialized", "dataset": "adsb_emergency.gold" }]
}
]
}
The order is resolved from dependencies, not from list position. A step can read an earlier step, or it can read a dataset another step writes. If a step materializes an Iceberg table, the next engine can pick it up there. If a step emits an ephemeral output, it stays inside the same engine and exists only for the run.
This distinction is useful in the UI as well. Airflow still sees one task for the pipeline. The console draws the step graph and the datasets between them, because that is the shape the operator actually needs when a run fails halfway through.
8. Monitoring Checks Need Their Own Model
One of the more useful lessons came from monitoring, not from pipeline execution.
The first freshness table could answer "what was the last verdict for this feed?" That is not the same as "what is the platform configured to watch?"
That distinction exposed several problems:
- stale rows for checks that no longer existed
- checks still pointing at paused DAGs after pipelines replaced them
- disabled checks that were indistinguishable from missing checks
- a notification cooldown that lived only in process memory and reset every cycle
The fix was to add two more tables.
monitoring_checks records what each agent is configured to watch, including disabled checks. alert_notifications records what was actually mailed.
Now the console can show four states:
okfailingdisablednot reporting
The last one is the important addition. A dead monitoring agent should not leave a calm-looking board just because nothing is writing new failures. Absence is a different fact from success, and it needs its own state.
The monitoring board groups checks by where the fault travels: stream, raw, silver, gold, and orchestration. It shows lag against each check's threshold, not as an absolute number. Four hours is fine for a daily table and an outage for an hourly one.
This is the same principle as the pipeline model: store the thing the operator will ask about later. A verdict alone was not enough. The configured check and the notification history were part of the operational truth.
9. What Changed Operationally
The immediate win was replacing hand-written DAGs without losing Airflow's scheduler.
Four hand-written DAGs were replaced by stored definitions:
| Removed DAG | Replacement definition |
|---|---|
duckdb_loader_readsb_hourly | pipeline_adsb_silver_hourly |
duckdb_loader_flightaware_hourly | pipeline_flightaware_silver_hourly |
adsb_curator_hourly | pipeline_adsb_gold_hourly |
flightaware_curator_hourly | pipeline_flightaware_gold_hourly |
The loaders were paused on 2026-08-14. The curators were paused on 2026-08-19. All four were deleted on 2026-08-20, after their replacements had run cleanly on schedule.
The cutover was intentionally boring. The replacement definitions were first run against live tables and checked for idempotence: same hour, same target, no duplicated rows, no missing rows, and row counts matching the path they replaced. For ADS-B silver, a replace-window run deleted and wrote the same bounded hour rather than appending a duplicate hour. For gold curation, rerunning an already-curated hour changed nothing, which is the expected behavior of the MERGE.
That cutover forced the surrounding system to improve:
- monitoring checks had to follow the new pipeline DAGs
- dataset catalog entries had to identify which tables were managed
- run records had to distinguish a counted
0from an unknown row count - previews had to sample raw prefixes in a way that was honest about what was read
- schema changes had to be visible in the run, not only in the table after the fact
The point is not that the new system is more elaborate. The point is that the platform now has a place to put these facts.
A DAG file can run a job. A pipeline definition can explain a job.
10. What I Would Do Differently
The model is useful because it was built after enough production friction to know what had to move out of DAG files. If I were starting it again, I would change a few things.
Build the dataset catalog first. The pipeline model depends on logical dataset references. Starting with inline bucket paths would have made the first pipeline faster to write and more expensive to repair.
Make schema behavior executable from day one. Having on_schema_change in the definition is not the same as enforcing it. A field that nobody reads is worse than no field, because it teaches operators to trust a promise the runtime is not keeping.
Avoid parallel run-history tables. The platform briefly had overlapping concepts for job runs and definition executions. Merging them into one pipeline_runs store made the console, monitoring, and lineage questions simpler because there was only one answer to "what ran?"
Treat monitoring configuration as data earlier. A verdict table is not enough. The platform also needs to know what each agent is configured to watch, which checks are disabled on purpose, and which alerts were actually sent.
11. Where This Goes Next
The next step is deeper data quality, not a separate orchestration system.
A quality check has the same shape as a pipeline: sources, a transform, a schedule, and a run record. The difference is that it asserts rather than writes a product table. The first version of that idea is already running as a raw-vs-silver parity check; the broader work is making those checks richer and more aviation-specific.
For example, raw-vs-silver parity for one ADS-B hour is naturally a check:
sources: adsb.raw, adsb.silver
transform: SQL producing observed_count, expected_count, passed
target: none
schedule: hourly
That does not need a separate quality subsystem on day one. It needs kind: "check" and a place in the monitoring surface.
Longer term, the same model supports more aviation-specific products: receiver coverage, reconstructed movement confidence, weather-to-traffic correlation, and operational checks that say whether a data product is fresh enough to use.
The broader lesson is the same one the lakehouse taught at the storage layer. Keep raw data replayable. Keep derived data versioned. Keep operational facts queryable.
Airflow still runs the work. The metadata tells us what the work means.