Skip to content

Dag processor kills callbacks for DAGs inside zip files at every bundle refresh #73858

Description

@TillerBurr

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:

  1. Put a DAG with a DAG-level on_success_callback that sleeps about 60 s into a zip in the dags folder.
  2. Keep the default bundle refresh interval, so the refresh runs while the callback sleeps.
  3. Trigger the DAG and let it succeed.
  4. 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:

  1. A DAG run succeeds and the scheduler creates a DagCallbackRequest for on_success_callback.
  2. 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.
  3. 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.
  4. 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?

  • Yes I am willing to submit a PR!

Code of Conduct

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    kind:bugThis is a clearly a bugneeds-triagelabel for new issues that we didn't triage yet

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions