Skip to content
6 changes: 3 additions & 3 deletions airflow/cli/commands/celery_command.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@
from airflow import settings
from airflow.configuration import conf
from airflow.utils import cli as cli_utils
from airflow.utils.cli import setup_locations, setup_logging
from airflow.utils.cli import setup_locations, setup_logging, DaemonContextWrapper
from airflow.utils.providers_configuration_loader import providers_configuration_loaded
from airflow.utils.serve_logs import serve_logs

Expand Down Expand Up @@ -80,7 +80,7 @@ def flower(args):
stdout.truncate(0)
stderr.truncate(0)

ctx = daemon.DaemonContext(
ctx = DaemonContextWrapper(
pidfile=TimeoutPIDLockFile(pidfile, -1),
stdout=stdout,
stderr=stderr,
Expand Down Expand Up @@ -227,7 +227,7 @@ def worker(args):
stdout_handle.truncate(0)
stderr_handle.truncate(0)

daemon_context = daemon.DaemonContext(
daemon_context = DaemonContextWrapper(
files_preserve=[handle],
umask=int(umask, 8),
stdout=stdout_handle,
Expand Down
5 changes: 3 additions & 2 deletions airflow/cli/commands/scheduler_command.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,8 @@
from airflow.jobs.job import Job, run_job
from airflow.jobs.scheduler_job_runner import SchedulerJobRunner
from airflow.utils import cli as cli_utils
from airflow.utils.cli import process_subdir, setup_locations, setup_logging, sigint_handler, sigquit_handler
from airflow.utils.cli import process_subdir, setup_locations, setup_logging, sigint_handler, \
sigquit_handler, DaemonContextWrapper
from airflow.utils.providers_configuration_loader import providers_configuration_loaded
from airflow.utils.scheduler_health import serve_health_check

Expand Down Expand Up @@ -69,7 +70,7 @@ def scheduler(args):
stdout_handle.truncate(0)
stderr_handle.truncate(0)

ctx = daemon.DaemonContext(
ctx = DaemonContextWrapper(
pidfile=TimeoutPIDLockFile(pid, -1),
files_preserve=[handle],
stdout=stdout_handle,
Expand Down
5 changes: 3 additions & 2 deletions airflow/cli/commands/triggerer_command.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@
from airflow.jobs.job import Job, run_job
from airflow.jobs.triggerer_job_runner import TriggererJobRunner
from airflow.utils import cli as cli_utils
from airflow.utils.cli import setup_locations, setup_logging, sigint_handler, sigquit_handler
from airflow.utils.cli import setup_locations, setup_logging, sigint_handler, sigquit_handler, DaemonContextWrapper
from airflow.utils.providers_configuration_loader import providers_configuration_loaded
from airflow.utils.serve_logs import serve_logs

Expand Down Expand Up @@ -68,7 +68,7 @@ def triggerer(args):
stdout_handle.truncate(0)
stderr_handle.truncate(0)

daemon_context = daemon.DaemonContext(
daemon_context = DaemonContextWrapper(
pidfile=TimeoutPIDLockFile(pid, -1),
files_preserve=[handle],
stdout=stdout_handle,
Expand All @@ -79,6 +79,7 @@ def triggerer(args):
triggerer_job_runner = TriggererJobRunner(
job=Job(heartrate=triggerer_heartrate), capacity=args.capacity
)

run_job(job=triggerer_job_runner.job, execute_callable=triggerer_job_runner._execute)
else:
signal.signal(signal.SIGINT, sigint_handler)
Expand Down
12 changes: 12 additions & 0 deletions airflow/utils/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
from pathlib import Path
from typing import TYPE_CHECKING, Callable, TypeVar, cast

import daemon
import re2
from sqlalchemy import select

Expand Down Expand Up @@ -381,3 +382,14 @@ def _wrapper(*args, **kwargs):
logging.disable(logging.NOTSET)

return cast(T, _wrapper)


class DaemonContextWrapper(daemon.DaemonContext):
"""Wrapper around DaemonContext to reset Stats instance in daemon process."""

def __enter__(self):
# in daemon context stats client needs to be reinitialized.
from airflow.stats import Stats
Stats.instance = None

super().__enter__()