Repository navigation
Fix scheduler crash when the Kafka event producer plugin is enabled - #72138
FrankYang0529 wants to merge 1 commit into
Conversation
Signed-off-by: PoAn Yang <payang@apache.org>
|
Hello @FrankYang0529 - thank you for your contributions to Apache Airflow! The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 30 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:
These pull requests stay open because maintainers are already engaged in them - they count towards your limit:
This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are. What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack. Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst. Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting |
Why
[kafka_event_producer] dag_run_events_enabled = True, the scheduler crashes withsqlalchemy.orm.exc.DetachedInstanceErrorwhen more than one Dag run is queued.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._start_queued_dagrunspasses. the second crashes ondag_run.dag_id.How
_build_producer, which runs on a single-useThreadPoolExecutor.settings.Sessiongives 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/unitbreeze shell --integration kafka --backend postgres --db-reset, thenpytest providers/apache/kafka/tests/integration/apache/kafka/plugins/test_event_producer.py --integration kafka.Was generative AI tooling used to co-author this PR?
{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.