Lesson 20 / 25
Concurrency, Pools and Priorities
Limit parallelism to protect shared systems and prioritise important work.
Do not stampede the database
Airflow can run many tasks at once, which can overwhelm a source database or API. Control it at several levels: max_active_runs (concurrent runs of one DAG), max_active_tasks (concurrent tasks within a DAG), parallelism (global limit), and pools, named sets of slots. Assign all tasks that hit one fragile system to a pool with, say, 3 slots; extra tasks wait in the queue. priority_weight lets important tasks go first when slots are scarce.
A pool in a task, and the idea in numbers
The task code assigns a pool. The small model that follows (run here) shows a pool of 3 slots with 3 running and 2 queued.
# in the DAG:
# BashOperator(task_id="export_a", bash_command="...", pool="source_db", priority_weight=5)
pool = 3; running = ["a", "b", "c"]; queued = ["d", "e"]
print(f"running {len(running)}/{pool}, queued {len(queued)}")
Output:
running 3/3, queued 2
Create pools in code review
Treat pool names and sizes as part of the design: document which system each protects, and review changes to them like any other configuration.
Quick check: What does a pool do?
- Stores XComs
- Limits how many tasks can use a shared resource at once
- Encrypts connections
- Renames DAGs
Answer
Limits how many tasks can use a shared resource at once — Pools cap concurrent use of a limited resource across DAGs.