Skip to content

Celery workers shutdown with kafka broker #73949

Description

@stephen-bracken

Under which category would you file this issue?

Providers

Apache Airflow version

3.3.1

What happened and how to reproduce it?

When running celery workers on kubernetes with a kafka broker, the Airflow celery app shuts down with a SIGKILL / 137 after a period of inactivity.

It also fails liveness probes with the following error:

Liveness probe failed: {"timestamp":"2026-09-30T09:29:48.369227Z","level":"debug","event":"Setting up DB connection pool (PID 643)","logger":"airflow.settings","filename":"settings.py","lineno":461} {"timestamp":"2026-09-30T09:29:48.369735Z","level":"debug","event":"settings.prepare_engine_args(): Using pool settings. pool_size=5, max_overflow=10, pool_recycle=1800, pid=643","logger":"airflow.settings","filename":"settings.py","lineno":585} {"timestamp":"2026-09-30T09:29:48.451782Z","level":"debug","event":"Could not retrieve value from section core, for key asset_manager_kwargs. Skipping redaction of this conf.","logger":"airflow.configuration","filename":"configuration.py","lineno":448} {"timestamp":"2026-09-30T09:29:49.068723Z","level":"debug","event":"registered serializers=[builtins.frozenset, builtins.set, builtins.tuple, datetime.date, datetime.datetime, datetime.timedelta, decimal.Decimal, deltalake.table.DeltaTable, kubernetes.client.models.v1_pod.V1Pod, kubernetes.client.models.v1_resource_requirements.V1ResourceRequirements, numpy.bool, numpy.bool_, numpy.complex128, numpy.complex64, numpy.float16, numpy.float32, numpy.float64, numpy.int16, numpy.int32, numpy.int64, numpy.int8, numpy.uint16, numpy.uint32, numpy.uint64, numpy.uint8, pandas.DataFrame, pandas.core.frame.DataFrame, pendulum.date.Date, pendulum.datetime.DateTime, pendulum.tz.timezone.FixedTimezone, pendulum.tz.timezone.Timezone, pydantic.main.BaseModel, pyiceberg.table.Table, uuid.UUID, zoneinfo.ZoneInfo] deserializers=[builtins.frozenset, builtins.set, builtins.tuple, datetime.date, datetime.datetime, datetime.timedelta, decimal.Decimal, deltalake.table.DeltaTable, numpy.bool, numpy.bool_, numpy.complex128, numpy.complex64, numpy.float16, numpy.float32, numpy.float64, numpy.int16, numpy.int32, numpy.int64, numpy.int8, numpy.uint16, numpy.uint32, numpy.uint64, numpy.uint8, pandas.DataFrame, pandas.core.frame.DataFrame, pendulum.date.Date, pendulum.datetime.DateTime, pendulum.tz.timezone.FixedTimezone, pendulum.tz.timezone.Timezone, pydantic.main.BaseModel, pyiceberg.table.Table, uuid.UUID, zoneinfo.ZoneInfo] stringifiers=[builtins.frozenset, builtins.set, builtins.tuple, deltalake.table.DeltaTable, pyiceberg.table.Table] in 3.607 ms","logger":"airflow.sdk.serde","filename":"__init__.py","lineno":509} {"timestamp":"2026-09-30T09:29:49.075477Z","level":"debug","event":"Could not retrieve value from section celery_result_backend_transport_options, for key sentinel_kwargs. Skipping redaction of this conf.","logger":"airflow.configuration","filename":"configuration.py","lineno":448} {"timestamp":"2026-09-30T09:29:49.278027Z","level":"debug","event":"Could not retrieve value from section celery, for key result_backend. Skipping redaction of this conf.","logger":"airflow.configuration","filename":"configuration.py","lineno":448} {"timestamp":"2026-09-30T09:29:49.328458Z","level":"debug","event":"Could not retrieve value from section database, for key sql_alchemy_engine_args. Skipping redaction of this conf.","logger":"airflow.configuration","filename":"configuration.py","lineno":448} {"timestamp":"2026-09-30T09:29:49.378452Z","level":"warning","event":"Skipping masking for a secret as it's too short (<5 chars)","logger":"airflow._shared.secrets_masker.secrets_masker","filename":"secrets_masker.py","lineno":605} {"timestamp":"2026-09-30T09:29:49.378565Z","level":"warning","event":"Skipping masking for a secret as it's too short (<5 chars)","logger":"airflow.sdk._shared.secrets_masker.secrets_masker","filename":"secrets_masker.py","lineno":605} {"timestamp":"2026-09-30T09:29:49.379892Z","level":"debug","event":"Could not retrieve value from section database, for key sql_alchemy_conn_async. Skipping redaction of this conf.","logger":"airflow.configuration","filename":"configuration.py","lineno":448} {"timestamp":"2026-09-30T09:29:49.414934Z","level":"debug","event":"Adding <function default_action_log at 0x7252b1fc47c0> to pre execution callback","logger":"airflow.utils.cli_action_loggers","filename":"cli_action_loggers.py","lineno":51} {"timestamp":"2026-09-30T09:29:49.568542Z","level":"debug","event":"Value for celery result_backend not found. Using sql_alchemy_conn with db+ prefix.","logger":"airflow.providers.celery.executors.default_celery","filename":"default_celery.py","lineno":137} {"timestamp":"2026-09-30T09:29:49.837849Z","level":"error","event":"KafkaError{code=TOPIC_AUTHORIZATION_FAILED,val=29,str=\"Subscribed topic not available: c0afee7b-9f07-3d5c-83a4-425135b95cb0.reply.celery.pidbox: Broker: Topic authorization failed\"}","logger":"kombu.transport.confluentkafka","filename":"confluentkafka.py","lineno":234} {"timestamp":"2026-09-30T09:29:55.841421Z","level":"debug","event":"Disposing DB connection pool (PID 643)","logger":"airflow.settings","filename":"settings.py","lineno":622} %4|1790760589.783|CONFWARN|rdkafka#consumer-3| [thrd:app]: Configuration property compression.codec is a producer property and will be ignored by this consumer instance Error: No nodes replied within time constraint

What you think should happen instead?

We are experimenting with using the kafka broker as a way to provide region specific celery queues for airflow tasks, so any guidance would be appreciated here.

Operating System

Debian GNU/Linux 12 (bookworm)

Deployment

Other Docker-based deployment

Apache Airflow Provider(s)

celery

Versions of Apache Airflow Providers

apache-airflow-providers-apache-kafka==1.15.1
apache-airflow-providers-celery==3.24.0
celery==5.6.3
confluent-kafka==2.15.0
kombu==5.6.2

Official Helm Chart version

No response

Kubernetes Version

1.33.5

Helm Chart configuration

No response

Docker Image customizations

      livenessProbe:
        exec:
          command:
            - sh
            - '-c'
            - >
              CONNECTION_CHECK_MAX_COUNT=0 exec /entrypoint python -m celery
              --app airflow.providers.celery.executors.celery_executor.app
              inspect ping -d celery@$(hostname)

Anything else?

Celery configuration:

{'accept_content': ['json'], 'event_serializer': 'json', 'worker_prefetch_multiplier': 1, 'task_acks_late': True, 'task_default_queue': 'airflow.tasks.default', 'task_default_exchange': 'airflow.tasks.default', 'task_track_started': True, 'broker_url': 'confluentkafka://< broker 1 >:9093;confluentkafka://< broker 2 >:9093;confluentkafka://< broker 3 >:9093;confluentkafka://< broker 4 >:9093', 'broker_transport_options': {'security_protocol': 'PLAINTEXT', 'sentinel_kwargs': {}, 'kafka_common_config': {'compression.type': 'zstd'}}, 'broker_connection_retry_on_startup': True, 'result_backend': 'db+postgresql+psycopg2://< airflow db >', 'result_backend_transport_options': {}, 'database_engine_options': {}, 'worker_concurrency': 16, 'worker_enable_remote_control': True, 'worker_redirect_stdouts': False, 'worker_hijack_root_logger': False}

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