Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 22 additions & 5 deletions airflow/bin/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -1151,9 +1151,9 @@ def worker(args):
"""Starts Airflow Celery worker"""
env = os.environ.copy()
env['AIRFLOW_HOME'] = settings.AIRFLOW_HOME
log = LoggingMixin().log

if not settings.validate_session():
log = LoggingMixin().log
log.error("Worker exiting... database connection precheck failed! ")
sys.exit(1)

Expand All @@ -1165,6 +1165,8 @@ def worker(args):
if autoscale is None and conf.has_option("celery", "worker_autoscale"):
autoscale = conf.get("celery", "worker_autoscale")
worker = worker.worker(app=celery_app) # pylint: disable=redefined-outer-name
skip_serve_logs = args.skip_serve_logs

options = {
'optimization': 'fair',
'O': 'fair',
Expand All @@ -1175,6 +1177,12 @@ def worker(args):
'loglevel': conf.get('core', 'LOGGING_LEVEL'),
}

if skip_serve_logs is False:
log.warning(
"Starting serve logs process within worker is going to be deprecated, "
"and will be removed from Airflow 2.0",
)

if conf.has_option("celery", "pool"):
options["pool"] = conf.get("celery", "pool")

Expand All @@ -1195,19 +1203,22 @@ def worker(args):
stderr=stderr,
)
with ctx:
sub_proc = subprocess.Popen(['airflow', 'serve_logs'], env=env, close_fds=True)
if skip_serve_logs is False:
sub_proc = subprocess.Popen(['airflow', 'serve_logs'], env=env, close_fds=True)
worker.run(**options)
sub_proc.kill()

stdout.close()
stderr.close()
else:
signal.signal(signal.SIGINT, sigint_handler)
signal.signal(signal.SIGTERM, sigint_handler)

sub_proc = subprocess.Popen(['airflow', 'serve_logs'], env=env, close_fds=True)
if skip_serve_logs is False:
sub_proc = subprocess.Popen(['airflow', 'serve_logs'], env=env, close_fds=True)

worker.run(**options)

if skip_serve_logs is False:
sub_proc.kill()


Expand Down Expand Up @@ -2242,6 +2253,12 @@ class CLIFactory:
'autoscale': Arg(
('-a', '--autoscale'),
help="Minimum and Maximum number of worker to autoscale"),
'skip_serve_logs': Arg(
("-s", "--skip_serve_logs"),
default=False,
help=(
"Don't start the serve logs process along with the workers."),
action="store_true"),
}
subparsers = (
{
Expand Down Expand Up @@ -2525,7 +2542,7 @@ class CLIFactory:
'func': worker,
'help': "Start a Celery worker node",
'args': ('do_pickle', 'queues', 'concurrency', 'celery_hostname',
'pid', 'daemon', 'stdout', 'stderr', 'log_file', 'autoscale'),
'pid', 'daemon', 'stdout', 'stderr', 'log_file', 'autoscale', 'skip_serve_logs'),
}, {
'func': flower,
'help': "Start a Celery Flower",
Expand Down
26 changes: 26 additions & 0 deletions tests/cli/test_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,32 @@ def test_ready_prefix_on_cmdline_dead_process(self):
with patch('psutil.Process', return_value=self.process):
self.assertEqual(get_num_ready_workers_running(self.gunicorn_master_proc), 0)

@mock.patch('airflow.bin.cli.subprocess.Popen')
def test_serve_logs_on_worker_start(self, mock_popen):
mock_popen.return_value.communicate.return_value = (b'output', b'error')
mock_popen.return_value.returncode = 0
args = self.parser.parse_args(['worker', '-c', '-1'])

with patch('celery.platforms.check_privileges') as mock_privil:
mock_privil.return_value = 0
with patch('celery.worker.WorkController.Blueprint.start') as mock_run:
mock_run.return_value.returncode = 0
cli.worker(args)
mock_popen.assert_called_once()

@mock.patch('airflow.bin.cli.subprocess.Popen')
def test_skip_serve_logs_on_worker_start(self, mock_popen):
mock_popen.return_value.communicate.return_value = (b'output', b'error')
mock_popen.return_value.returncode = 0
args = self.parser.parse_args(['worker', '-c', '-1', '-s'])

with patch('celery.platforms.check_privileges') as mock_privil:
mock_privil.return_value = 0
with patch('celery.worker.WorkController.Blueprint.start') as mock_run:
mock_run.return_value.returncode = 0
cli.worker(args)
mock_popen.assert_not_called()

def test_cli_webserver_debug(self):
env = os.environ.copy()
proc = psutil.Popen(["airflow", "webserver", "-d"], env=env)
Expand Down