GitHub user zhanqian-zhang2 created a discussion: Task memory diagnostics in
Airflow 3: measurement scopes and trade-offs
## Context and goal
Airflow 2.10 introduced periodic task CPU and memory observations in #39650.
Those metrics are no longer emitted in Airflow 3; their absence was reported in
#49983. #51602 later removed stale documentation, and #56690 explored restoring
CPU and memory metrics but was not merged.
Before attempting another implementation, I would like to clarify **what a task
memory metric should promise to measure**.
The practical motivation is:
> Help users notice tasks whose memory behavior may be worth investigating or
> optimizing.
There are several useful ways to answer that question, with different
measurement boundaries and implementation costs.
The scopes below are **design alternatives / possible extensions, not an
implementation commitment or a fixed roadmap**.
In particular, the smallest scope could provide useful evidence when high
memory is actually observed, but it would not reliably rank executions by
historical peak.
## Scope 1 — Python Task Runner RSS diagnostic
The smallest useful signal would retain the lightweight periodic
process-sampling idea from Airflow 2.10 while defining its meaning more
narrowly.
An illustrative metric name is:
```
task.runner_rss_bytes
```
with the contract:
> Periodically sampled Resident Set Size of one Python Task Runner process,
> reported in bytes.
### Measurement boundary
Conceptually:
```
Supervisor
|
+-- Python Task Runner <-- measured
|
+-- task/operator code
+-- possible subprocesses
```
Only the Python Task Runner PID is measured.
It would **not** mean:
- total memory consumed by the execution workload;
- memory consumed by subprocesses;
- process-tree memory;
- container or Pod memory;
- host memory pressure;
- the total or maximum across concurrent executions of the same logical task;
- a guaranteed execution peak.
This distinction matters because some tasks perform most of their work in Bash,
another Python interpreter, Java, ffmpeg, or another local/remote process.
### Sampling placement and lifecycle
I do not think the exact sampling location should be decided in the metric
contract.
Airflow 2.10 observed the task subprocess externally from its task runner. In
Airflow 3, the Supervisor starts and tracks the Python Task Runner process, so
Supervisor-side observation is one possible design.
Sampling from inside the Python runner is another possibility and has different
access to task identity and initialized metrics.
Before choosing either, the implementation should validate:
- when sufficient task identity/tags become available;
- startup and finalization coverage;
- `run_as_user` / re-exec behavior;
- process permissions;
- shutdown/export behavior;
- thread/fork safety.
A realistic lifecycle contract is therefore:
> Start observation once sufficient task identity and measurement context are
> available, and continue through finalization where feasible. Define the exact
> boundaries in the implementation.
Observation failures must remain strictly best-effort and must never affect
task result, retry, deferral, or cleanup.
### What Scope 1 can and cannot answer
For Python-heavy tasks, runner RSS may still be a useful low-cost signal:
> "This Python runner was observed using substantial memory."
It cannot reliably answer:
> "These are the tasks with the highest historical execution peaks."
For example:
```
0.0s 500 MiB
1.0s 8 GiB
2.0s 600 MiB
3.0s exit
```
A periodic sampler may never observe the 8 GiB spike, and a sufficiently short
execution may finish before a useful observation reaches the backend.
That is an inherent limitation of sampled current-state diagnostics.
## The main unresolved Scope-1 problem: producer identity and concurrency
Consider two concurrent executions of the same logical task:
```
dag_id=foo
task_id=transform
execution A runner RSS = 900 MiB
execution B runner RSS = 200 MiB
```
Useful logical-task values could be:
```
current total = 1100 MiB
current max = 900 MiB
active count = 2
```
But two independent producers emitting individual Gauge values do not
automatically produce any of those meanings.
`(dag_id, task_id)` is a useful **business grouping key**, but it is not
necessarily sufficient producer/stream identity.
If independent runners appear as the same telemetry stream, their observations
may be conflated.
The obvious alternative—adding `run_id`, PID, execution UUID, or another
execution identity—distinguishes producers but introduces execution-level
time-series cardinality and churn.
Different metrics backends expose this problem differently. That is an
implementation concern, but the underlying semantic problem is
backend-independent.
So Scope 1 has an important prerequisite:
> Can Airflow export independent runner observations with unambiguous and
> reasonably bounded identity?
If not, the apparently simple runner metric may not be a sound abstraction, and
aggregation before export may be preferable.
I would prefer to decide this semantic question first and then determine the
appropriate StatsD/OpenTelemetry implementation, rather than letting a
telemetry primitive define the metric contract.
## Scope 2 — aggregate current runner observations
A stronger design answers:
> "How much runner memory are the currently active executions of this logical
> task using?"
This requires aggregation before the low-cardinality logical-task metric is
exported.
Conceptually:
```
execution A --+
|
execution B --+--> stable local aggregation owner
| |
execution C --+ +-- group by (dag_id, task_id)
|
+-- SUM current runner RSS
+-- MAX current runner RSS
+-- COUNT active executions
```
The important property is that the aggregation owner has a lifetime longer than
one task execution and can maintain bounded state for currently active
executions.
Execution-level information can exist internally:
```
execution identity
PID / process identity
dag_id
task_id
```
without becoming permanent metric dimensions.
This separates:
```
high-cardinality execution identity
-> bounded local state
logical task identity
-> exported metrics
```
### Possible metric semantics
For each logical task within one aggregation owner:
| Illustrative signal | Meaning
|
| ---------------------------------- |
--------------------------------------------------------- |
| `task.active_runner_rss_bytes` | Sum of successfully observed active
runner RSS values |
| `task.max_active_runner_rss_bytes` | Maximum of successfully observed active
runner RSS values |
| `task.active_executions` | Executions known by the owner to be
active |
| `task.observed_executions` | Active executions whose RSS was
successfully read |
For example, if three executions are alive but only two can be sampled:
```
active_executions = 3
observed_executions = 2
```
The RSS total/max then describe the **observed subset**.
Reporting `active_executions = 2` would incorrectly redefine "active" as
"successfully observed".
Also, even here, summing process RSS should not be interpreted as deduplicated
physical memory.
### Multiple aggregation owners
"Local" does not imply deployment-wide aggregation.
Suppose:
```
owner A:
total = 1100 MiB
max = 900 MiB
count = 2
owner B:
total = 500 MiB
max = 500 MiB
count = 1
```
A global logical-task view would conceptually require:
```
total -> SUM = 1600 MiB
max -> MAX = 900 MiB
count -> SUM = 3
```
The sources must therefore remain distinguishable until that second aggregation
occurs.
This cross-owner identity/query contract is part of Scope 2; it is not implied
merely by choosing metric names.
### Lifecycle complexity
The owner would need explicit handling for:
- execution registration/unregistration;
- retry;
- defer/resume;
- unexpected process termination;
- Supervisor/worker shutdown;
- stale registrations;
- PID reuse;
- collection racing with execution termination.
PID alone may be insufficient as durable process identity because PIDs can be
reused.
### Executor implications
The semantic contract need not be executor-specific, but the concrete
aggregation owner may be.
Possible candidates include:
- **CeleryExecutor:** a long-lived Celery worker domain is a plausible first
aggregation boundary;
- **LocalExecutor:** a long-lived parent/management layer may provide an
equivalent boundary, but this needs code-level validation;
- **KubernetesExecutor:** task-per-Pod execution has no equivalent
cross-execution owner inside each task Pod, so aggregation would require a
different architecture.
These are design candidates, not claims that Airflow already provides such
aggregation facilities.
A practical implementation could initially support only execution paths for
which a correct stable owner has been identified.
An unsupported execution path should not silently emit runner-local readings
under the same metric name as a metric promising logical-task aggregation.
## Scope 3 — execution-workload memory accounting
Scopes 1 and 2 still measure **Python Task Runner processes**.
That may substantially underestimate executions such as:
```
Task Runner 300 MiB
|
+-- Python 700 MiB
+-- Java 1.5 GiB
+-- ffmpeg 2.0 GiB
```
If the intended question is:
> "How much local memory did this execution workload consume?"
then the measurement boundary itself needs to change.
### Process discovery as an approximation
One possible approach is to discover processes associated with an execution
through parent/child relationships or a process-group boundary where one exists.
This could improve coverage of common local subprocesses, but it should be
treated as:
> a process-discovery approach with incomplete attribution,
not as a ready-made resource-accounting boundary.
Limitations include:
- processes may detach or escape the discovered hierarchy/group;
- other coordinators/runtimes may use different process models;
- summing RSS can double-count shared resident pages;
- PSS has different attribution semantics and higher collection
cost/availability constraints.
Any metric based on this approach should state exactly which processes and
memory measure it represents.
### Dedicated per-execution cgroup
A stronger Linux-specific design would establish a dedicated cgroup for an
execution **before the workload starts**:
```
worker
|
+-- execution A cgroup
| +-- Task Runner
| +-- bash
| +-- ffmpeg
|
+-- execution B cgroup
+-- Task Runner
+-- java
```
The kernel could then account memory to that execution resource group and its
descendants.
Relevant cgroup-v2 information to investigate includes:
```
memory.current
memory.stat
memory.peak
memory.events
```
A working-set-style value could also be considered if that is the desired
semantic.
The important improvement is not that a cgroup creates one universally
"correct" memory value—shared memory, cache and charge ownership still have
semantics that need documentation.
The improvement is the **measurement boundary**:
> memory is accounted to the execution workload rather than inferred from one
> selected process.
### Historical execution peaks
A kernel-maintained value such as `memory.peak` may also capture peaks that
periodic polling misses.
That could support a separate concept such as:
```
task.execution_peak_memory_bytes
```
with one peak observation for each execution whose final resource state was
successfully captured.
This is different from periodically recording RSS samples into a Histogram:
```
periodic RSS Histogram
-> distribution of sampling points
per-execution peak distribution
-> distribution of execution peaks
```
The second is more useful for historical memory-sizing analysis.
However, this is not automatic.
A surviving lifecycle owner must read and export the peak **before the cgroup
is removed**.
In particular, a design that only records the peak after successful task
completion would miss some of the most interesting executions:
- failures;
- kills;
- OOM termination.
Reliable peak collection therefore requires an observer/lifecycle design that
survives long enough to handle abnormal termination.
It would also require an appropriate bytes-valued distribution API and backend
representation; a duration-oriented timing API should not be reused merely
because it already exposes a Histogram-like primitive.
### Platform cost
Per-execution cgroups introduce substantial platform questions:
- cgroup-v2 availability;
- delegation and permissions;
- workers running inside containers;
- ownership and cleanup;
- abnormal termination;
- unsupported operating systems/deployments.
If this scope is implemented, lack of cgroup support should not silently change
the same metric into runner RSS.
## The scopes should remain semantically distinct
These are different measurements:
```
one runner's RSS
aggregate current RSS of active runners
process-tree/process-group approximation
execution-cgroup memory
```
They answer different questions.
They should not all silently become implementations of:
```
task.memory_usage
```
depending on executor, platform, or available privileges.
Likewise, changing a published metric from:
```
one runner observation
```
to:
```
sum of all active runners
```
would materially change dashboard and alert semantics.
If more than one scope is useful, separate metric names or an explicit
migration strategy would be safer.
## Questions for the community
I would particularly appreciate feedback on these points:
1. **Is a precisely named Python Task Runner RSS diagnostic useful in Airflow
core at all?**
2. **If so, can independent concurrent runner observations be exported with
sufficiently clear and bounded producer identity?**
If not, should a stable local owner and Scope-2 aggregation be the minimum
design?
3. **Should a** `**(dag_id, task_id)**` **memory metric be expected to
represent correctly aggregated concurrent state from its first introduction?**
4. **Longer term, is there value in exposing both runner diagnostics and a
separate execution-workload memory signal?**
5. **If workload-level accounting is desirable, is an Airflow-managed
per-execution cgroup worth investigating, or should that responsibility remain
with the deployment/container layer?**
My current inclination is to keep any first contribution as small as possible,
but not at the cost of publishing a metric whose measurement boundary, producer
identity, or aggregation semantics are ambiguous.
GitHub link: https://github.com/apache/airflow/discussions/73645
----
This is an automatically sent email for [email protected].
To unsubscribe, please send an email to: [email protected]