Repository navigation
Task memory diagnostics in Airflow 3: measurement scopes and trade-offs #73623
Unanswered
zhanqian-zhang
asked this question in
Ideas
Replies: 0 comments
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
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:
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:
with the contract:
Measurement boundary
Conceptually:
Only the Python Task Runner PID is measured.
It would not mean:
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:
run_as_user/ re-exec behavior;A realistic lifecycle contract is therefore:
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:
It cannot reliably answer:
For example:
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:
Useful logical-task values could be:
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:
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:
This requires aggregation before the low-cardinality logical-task metric is exported.
Conceptually:
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:
without becoming permanent metric dimensions.
This separates:
Possible metric semantics
For each logical task within one aggregation owner:
task.active_runner_rss_bytestask.max_active_runner_rss_bytestask.active_executionstask.observed_executionsFor example, if three executions are alive but only two can be sampled:
The RSS total/max then describe the observed subset.
Reporting
active_executions = 2would 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:
A global logical-task view would conceptually require:
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:
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:
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:
If the intended question is:
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:
not as a ready-made resource-accounting boundary.
Limitations include:
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:
The kernel could then account memory to that execution resource group and its descendants.
Relevant cgroup-v2 information to investigate includes:
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:
Historical execution peaks
A kernel-maintained value such as
memory.peakmay also capture peaks that periodic polling misses.That could support a separate concept such as:
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:
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:
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:
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:
They answer different questions.
They should not all silently become implementations of:
depending on executor, platform, or available privileges.
Likewise, changing a published metric from:
to:
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:
If not, should a stable local owner and Scope-2 aggregation be the minimum design?
**(dag_id, task_id)**memory metric be expected to represent correctly aggregated concurrent state from its first introduction?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.
All reactions