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.