Under which category would you file this issue?
Airflow Core
Apache Airflow version
3.1.7 and current main
What happened and how to reproduce it?
When DAGs are shipped as a zip in the local dags-folder bundle, the dag processor sometimes SIGKILLs a DAG callback while it is still running. The callback never completes. No file changed and no deploy ran. The trigger is the periodic bundle refresh.
The kill only happens when a callback is still running at the moment a refresh runs.
This is related to #66483, but it's a different bug. #66483 was a mismatch on bundle_version, and #66484 fixed it by comparing presence_key = (bundle_name, rel_path). Here the mismatch is on rel_path itself, so presence_key doesn't match either, and main is still affected.
Reproduction:
- Put a DAG with a DAG-level on_success_callback that sleeps about 60 s into a zip in the dags folder.
- Keep the default bundle refresh interval, so the refresh runs while the callback sleeps.
- Trigger the DAG and let it succeed.
- Watch the dag processor log. You should see Stopping processor for zip/dag.py followed by exit_code=<Negsignal.SIGKILL: -9>, and the callback's side effect never happens.
The sequence we observed in production:
- A DAG run succeeds and the scheduler creates a DagCallbackRequest for on_success_callback.
- The dag processor queues the callback. Its DagFileInfo is keyed by the DAG's path inside the zip: rel_path=my_dags.zip/my_dag.py.
- About 11 s later the timed bundle refresh runs (LocalDagBundle.refresh() is a no-op, but the refresh still rescans). _find_files_in_bundle() returns top-level paths only, including my_dags.zip.
- terminate_orphan_processes() finds that my_dags.zip/my_dag.py isn't in the scanned set and kills the processor:
WARNING - Stopping processor for my_dags.zip/my_dag.py
INFO - Process exited pid=133594 exit_code=<Negsignal.SIGKILL: -9> signal_sent=SIGKILL
The callback, a Slack notification, was never sent.
Minimal example showing keys don't match:
from pathlib import Path
from unittest import mock
import airflow
from airflow.dag_processing.manager import DagFileInfo, DagFileProcessorManager
bundle = dict(bundle_name="dags-folder", bundle_path=Path("/opt/dags"))
callback_file = DagFileInfo(rel_path=Path("my_dags.zip/my_dag.py"), **bundle) # key from _add_callback_to_queue
scanned_file = DagFileInfo(rel_path=Path("my_dags.zip"), **bundle) # what the bundle refresh finds
print(callback_file.presence_key == scanned_file.prescence_key) # False
# Get killed as an orphan
manager = DagFileProcessorManager(max_runs=1)
processor = mock.MagicMock()
manager._processors = {callback_file: processor}
manager.terminate_orphan_processes(present={scanned_file})
print("airflow", airflow.__version__)
print("callback processor killed:", processor.kill.called, processor.kill.call_args)
What you think should happen instead?
A callback for a DAG inside a zip should not be treated as orphaned while its containing zip is still present in the bundle scan. File presence for my_dags.zip/my_dag.py should be satisfied by my_dags.zip being present.
Are you willing to submit PR?
Code of Conduct
Under which category would you file this issue?
Airflow Core
Apache Airflow version
3.1.7 and current main
What happened and how to reproduce it?
When DAGs are shipped as a zip in the local dags-folder bundle, the dag processor sometimes SIGKILLs a DAG callback while it is still running. The callback never completes. No file changed and no deploy ran. The trigger is the periodic bundle refresh.
The kill only happens when a callback is still running at the moment a refresh runs.
This is related to #66483, but it's a different bug. #66483 was a mismatch on
bundle_version, and #66484 fixed it by comparingpresence_key = (bundle_name, rel_path). Here the mismatch is onrel_pathitself, sopresence_keydoesn't match either, and main is still affected.Reproduction:
The sequence we observed in production:
WARNING - Stopping processor for my_dags.zip/my_dag.py
INFO - Process exited pid=133594 exit_code=<Negsignal.SIGKILL: -9> signal_sent=SIGKILL
The callback, a Slack notification, was never sent.
Minimal example showing keys don't match:
What you think should happen instead?
A callback for a DAG inside a zip should not be treated as orphaned while its containing zip is still present in the bundle scan. File presence for
my_dags.zip/my_dag.pyshould be satisfied bymy_dags.zipbeing present.Are you willing to submit PR?
Code of Conduct