Skip to content

Fix scheduler crash when the Kafka event producer plugin is enabled - #72138

Open
FrankYang0529 wants to merge 1 commit into
apache:mainfrom
FrankYang0529:airflow-kafka-event-producer-session-detach
Open

Fix scheduler crash when the Kafka event producer plugin is enabled#72138
FrankYang0529 wants to merge 1 commit into
apache:mainfrom
FrankYang0529:airflow-kafka-event-producer-session-detach

Conversation

@FrankYang0529

Copy link
Copy Markdown
Member

Why

  • With [kafka_event_producer] dag_run_events_enabled = True, the scheduler crashes with sqlalchemy.orm.exc.DetachedInstanceError when more than one Dag run is queued.
  • The DagRun listeners run inside the scheduler's own transaction. It moves queued Dag runs to RUNNING.
  • Building the Kafka producer resolves the Kafka connection. On the scheduler that means MetastoreBackend.get_connection, which is decorated with @provide_session. create_session() is thread-scoped, so on the scheduler's thread it reuses the scheduler's own session and closes it on exit.
  • The first Dag run in _start_queued_dagruns passes. the second crashes on dag_run.dag_id.

How

  • Producer construction moves into _build_producer, which runs on a single-use ThreadPoolExecutor. settings.Session gives out one session per thread, so the connection lookup gets a session of its own and leaves the caller's transaction alone.

Verification

  • uv run --project providers/apache/kafka pytest providers/apache/kafka/tests/unit
  • breeze shell --integration kafka --backend postgres --db-reset, then pytest providers/apache/kafka/tests/integration/apache/kafka/plugins/test_event_producer.py --integration kafka.

Was generative AI tooling used to co-author this PR?
  • Yes - Claude Code

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Signed-off-by: PoAn Yang <payang@apache.org>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant