diff --git a/providers/celery/src/airflow/providers/celery/cli/celery_command.py b/providers/celery/src/airflow/providers/celery/cli/celery_command.py index 987233241e105..3a953f2a26b8b 100644 --- a/providers/celery/src/airflow/providers/celery/cli/celery_command.py +++ b/providers/celery/src/airflow/providers/celery/cli/celery_command.py @@ -310,6 +310,9 @@ def worker(args): if not celery_log_level: celery_log_level = config.get("logging", "LOGGING_LEVEL") + if args.verbose: + celery_log_level = "DEBUG" + # Setup Celery worker options = [ "worker", diff --git a/providers/celery/tests/unit/celery/cli/test_celery_command.py b/providers/celery/tests/unit/celery/cli/test_celery_command.py index b189f0f9078ae..8a91aac7efc8a 100644 --- a/providers/celery/tests/unit/celery/cli/test_celery_command.py +++ b/providers/celery/tests/unit/celery/cli/test_celery_command.py @@ -233,6 +233,39 @@ def test_worker_applies_celery_mp_start_method( mock_set_mp.assert_called_once_with("celery") +@pytest.mark.usefixtures("conf_stale_bundle_cleanup_disabled") +class TestWorkerLogLevel: + @pytest.fixture(autouse=True) + def _disable_cli_action_logging(self): + with ( + patch("airflow.utils.cli.cli_action_loggers.on_pre_execution"), + patch("airflow.utils.cli.cli_action_loggers.on_post_execution"), + ): + yield + + @classmethod + def setup_class(cls): + with conf_vars({("core", "executor"): "CeleryExecutor"}): + importlib.reload(executor_loader) + importlib.reload(cli_parser) + cls.parser = cli_parser.get_parser() + + @conf_vars({("logging", "celery_logging_level"): "INFO"}) + @mock.patch("airflow.providers.celery.cli.celery_command.setup_locations") + @mock.patch("airflow.providers.celery.cli.celery_command.Process") + @mock.patch("airflow.providers.celery.executors.celery_executor.app") + def test_worker_verbose_overrides_configured_celery_loglevel( + self, mock_celery_app, mock_popen, mock_locations + ): + mock_locations.return_value = ("pid_file", None, None, None) + args = self.parser.parse_args(["celery", "worker", "--verbose", "--skip-serve-logs"]) + + celery_command.worker(args) + + worker_options = mock_celery_app.worker_main.call_args[0][0] + assert worker_options[worker_options.index("--loglevel") + 1] == "DEBUG" + + @pytest.mark.backend("mysql", "postgres") @pytest.mark.usefixtures("conf_stale_bundle_cleanup_disabled") class TestWorkerMultiTeam: