dag_processing.processes metric missing decrement on graceful process exit
- Dominant language
- Python
- Stars
- 46.9k
- Forks
- 17.8k
- Avg merge
- 2d 9h
- Merged PRs (30d)
- 472
Description
### Under which category would you file this issue?
Airflow Core
### Apache Airflow version
3.3.0
### What happened and how to reproduce it?
The dag_processing.processes metric description states: "the delta is negative when, since the last metric was sent, processes have completed." However, `_collect_results()` , the normal graceful completion path, emits no stats.decr, so the delta is never negative under normal operation. A decr is only emitted for abnormal exits , which are : stop (file deleted), timeout, and terminate (shutdown). This means the metric does not behave as documented.
**Current behaviour**
`stats.incr("dag_processing.processes", tags={"action": "start", ...})` is called
in `_start_new_processes()` when a child `DagFileProcessorProcess` is spawned.
A matching `stats.decr(...)` is called in three abnormal exit paths:
- `terminate_orphan_processes()` action: `"stop"` (file deleted mid-parse)
- `_kill_timed_out_processors()` action: `"timeout"` (exceeded `processor_timeout`)
- `terminate()` action: `"terminate"` (manager shutting down)
However, `_collect_results()` , the **normal graceful exit path** called every loop
iteration pops the finished processor from `self._processors` and calls
`processor.close()` with **no corresponding `decr`**:
This means under normal operation (parsers completing successfully), the metric value drifts upward indefinitely and never accurately reflects the number of currently running processes
### What you think should happen instead?
A `stats.decr("dag_processing.processes", tags={"action": "finish", ...})` should be emitted when a process completes gracefully, symmetric with the incr in `_start_new_processes()`.
### Operating System
_No response_
### Deployment
None
### Apache Airflow Provider(s)
_No response_
### Versions of Apache Airflow Providers
_No response_
### Official Helm Chart version
Not Applicable
### Kubernetes Version
_No response_
### Helm Chart configuration
_No response_
### Docker Image customizations
_No response_
### Anything else?
_No response_
### Are you willing to submit PR?
- [x] Yes I am willing to submit a PR!
### Code of Conduct
- [x] I agree to follow this project's [Code of Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
Contributor guide
Research direction
Start by reading _start_new_processes() and _collect_results(), then compare the existing decr calls in terminate_orphan_processes(), _kill_timed_out_processors(), and terminate(). Done means graceful processor completion emits a matching stats.decr with action "finish" and the relevant DAG processing metric tests pass.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- python
- Domain
- observability
- Issue type
- Bug
- Difficulty
- 3/5
- Estimated time
- 1-2 days
- Activity status
- Quiet
- Clarity
- Mostly clear
- Newbie friendliness
- 68/100