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]

Reply via email to