Week 7: Orchestration Across Projects with Apache Airflow#
Note
The Airflow chapter for this week is posted before class. Until then, this page is the agenda.
This week we step up one level in the automation stack.
doit orders the steps inside one project. A workflow orchestrator orders many
projects: it runs each pipeline on a schedule or when the data it depends on
lands, retries what fails, records every run, and can backfill a gap after an
outage. Apache Airflow is the most widely used of these tools, and the class
report is the right thing to point it at, because the report is not one pipeline.
It is twenty-odd small producer pipelines with a single consumer that should
rebuild when all of them have landed.
Announcements#
HW 4 is due tonight, at the start of class. It is the last of the four coding assignments.
HW 5 is due next week, Tuesday, November 17. Your pull request should be open by now, so your two reviewers have time. Tonight’s Airflow lab runs on the report’s pipelines. It is done in class and is not graded.
Tonight’s material is on the exam, Tuesday, December 1. See Exam Preparation.
Objectives#
Say what a workflow orchestrator does that a task runner does not, and name the cases where reaching for one is justified.
Read an Airflow DAG: tasks, dependencies, the schedule, and the data interval.
Explain what a data interval is and why a run for interval t must use only what was known at t. This is the same discipline the equity-premium literature imposes on an out-of-sample forecast, enforced by the scheduler instead of by hand.
Generate DAGs from a manifest rather than writing them by hand, so that twenty contributions do not mean twenty hand-maintained files.
Schedule a consumer on upstream assets, so the report DAG fires when its inputs have all updated, with AND semantics across many producers.
Diagnose a failure from the scheduler’s UI and its logs, distinguish a code bug from an expired credential from a non-idempotent task, and backfill after an outage.
Explain what makes a task idempotent and why a scheduler is unusable without it.
Say what containers and Kubernetes add beneath an orchestrator, and how a container image relates to the
uvandpixilockfiles from week 2.
Agenda Item 1: Why Another Scheduler?#
You already have doit, and cron is coming in week 8. Airflow is not a
replacement for either, and the honest version of this week starts with when not
to use it.
doitknows one project’s dependency graph and runs it now.cronruns a command at a time. Neither knows that project B’s input is project A’s output, neither records run history, and neither can tell you which of last Tuesday’s runs failed and why.What an orchestrator adds: a graph across projects, run history, retries with backoff, alerting, and backfill.
What it costs: a scheduler, a metadata database, and a webserver to keep alive. This is the trade every data platform team argues about, and you should be able to argue both sides.
Agenda Item 2: Airflow Concepts, Then the Course’s DAGs#
DAGs, tasks, and operators. Logical date versus wall-clock time, and the data interval.
catchupand what it means to have a scheduler decide it owes you fourteen runs.Assets (datasets): scheduling a DAG on the arrival of data rather than on a clock, and AND semantics when a consumer depends on many producers.
Generated DAGs. Every contribution already declares its outputs in a
chartbook.tomlmanifest. So its DAG is generated from the manifest, one producer DAG per contribution, each running that project’sdoitand then itsdoit test. Nobody hand-writes Airflow wiring, and a new contribution appears in the graph by virtue of having a manifest.The report DAG is then
schedule=[Asset(...), ...]over the contributions’ outputs: it rebuilds when they have all landed. The report becomes a document that regenerates itself whenever any input updates, which is the thesis of this course in one artifact.
Agenda Item 3: Operating It#
Reading about a scheduler teaches you nothing. You run one.
airflow standaloneagainst your clone, with a compressed clock so that a quarter’s worth of run history accumulates in minutes.Read the UI: the grid, a task’s logs, and the difference between a task that failed and a task that never ran because its upstream did.
The incident lab. The faults are the real ones: an upstream renames a column and twelve charts go red; a WRDS credential expires and fails in a way that looks nothing like a code bug; a contribution appends instead of overwrites and is therefore not idempotent, so a retry corrupts it.
Backfill after a simulated outage, and confirm the result matches what a timely run would have produced.
The shared instance stays a projector demo that I run. You run yours locally.
Agenda Item 4: Where the Tasks Run — Containers and Kubernetes#
Briefly, because the job postings name these tools and an orchestrator is where you meet them.
Containers. A Docker image freezes the whole environment a task runs in, operating system included. That is one level below the conda,
uv, andpixilockfiles from week 2, and the same image runs on your laptop, a CI runner, or a cluster.Kubernetes schedules containers across a fleet of machines. Production Airflow commonly runs each task in its own container on Kubernetes (the
KubernetesExecutor, or theKubernetesPodOperatorfor a single task), so twenty producer pipelines with conflicting dependencies can share one scheduler.We do not run a cluster. The goal is to recognize the layers when a posting names them: the lockfile pins the packages, the image pins the machine, and the orchestrator decides when and where the image runs.
Agenda Item 5: The Paper#
Goyal, Welch, and Zafirov (2024), A Comprehensive 2022 Look at the Empirical Performance of Equity Premium Prediction, extends the original seventeen predictors to forty-six and asks whether any of them beats the historical mean out of sample. The answer is sobering, and the discipline behind it is the reason this paper belongs in the orchestration week: a forecast made at date t may use only what was known at t. Doing that by hand is error-prone. A scheduler with correct data intervals makes it structural. Week 9 takes the dataset you are all building and actually fits models to it.
Looking ahead to Week 8#
Week 8 is CI/CD with GitHub Actions, and the contrast with what you ran tonight is the point of the week. Both are automation triggered by an event; one is triggered by a push and one by a clock or a dataset. The case study is a live FedWatch monitor that rebuilds itself every morning, computing the policy surprises of Bernanke and Kuttner (2005) from the futures data you already pulled in HW 2.