Skip to content

Handle concurrent XCom writes with savepoint and UPDATE fallback to avoid 409 - #69959

Open
nagasrisai wants to merge 15 commits into
apache:mainfrom
nagasrisai:fix/xcom-execution-api-concurrent-409
Open

nagasrisai wants to merge 15 commits into
apache:mainfrom
nagasrisai:fix/xcom-execution-api-concurrent-409

Conversation

@nagasrisai

@nagasrisai nagasrisai commented Jul 16, 2026 •

Copy link
Copy Markdown
Contributor

Closes #69956

Summary

When two task runners (concurrent execution or task retry) write the same XCom key at roughly the same time, both call XComModel.set which does a DELETE followed by an INSERT. Both DELETEs succeed (they each see the record that was just written by the peer or see nothing), then both try to INSERT the same primary key. The second INSERT hits the unique constraint and the Execution API converts the IntegrityError to an HTTP 409, causing the task to fail.

Fix

Wrap the XComModel.set call in a savepoint (session.begin_nested()) so that a constraint violation only rolls back the XCom write and not any earlier work in the same request (e.g. the task_map merge). The IntegrityError is caught and handled with an explicit UPDATE that overwrites the existing entry with the latest value, which is the correct semantic for a retry.

try:
    with session.begin_nested():
        XComModel.set(...)
except IntegrityError:
    # concurrent write already committed; overwrite with latest value
    session.execute(update(XComModel).where(...).values(value=value, dag_result=dag_result))

Changes

  • airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py
    • Add update to the sqlalchemy import
    • Add IntegrityError import from sqlalchemy.exc
    • Wrap XComModel.set in session.begin_nested() and catch IntegrityError with an UPDATE fallback

Important

🛠️ Maintainer triage note for @nagasrisai · by @potiuk · 2026-07-28 16:11 UTC

Helpful heads-up from the maintainers — please address before this PR can be reviewed:

  • ❌ Pre-commit / static checks. See docs.
  • ❌ Unresolved review comments: 2 thread(s). See docs.

The ball is in your court — you've been assigned to this PR. Fix the above, then mark it Ready for review.

See the Pull Request quality criteria for how to fix each item. There is no rush.

Note: your branch is 270 commits behind main — please rebase and push again to get up-to-date CI results.

Automated triage — may be imperfect; a maintainer takes the next look. We use this two-stage triage process so maintainers' limited time goes to the conversation with you.

nagasrisai added a commit to nagasrisai/airflow that referenced this pull request Jul 16, 2026
@potiuk

potiuk commented Jul 20, 2026

Copy link
Copy Markdown
Member

@nagasrisai This PR has a few issues that need to be addressed before it can be reviewed — please see our Pull Request quality criteria.

Issues found:

  • ❌ Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main --stage pre-commit locally to reproduce and fix.

Note: Your branch is 51 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • Fix each issue listed above.
  • Make sure static checks pass locally (prek run --from-ref main --stage pre-commit).
  • Mark the PR as "Ready for review" when you're done.

There is no rush — take your time and work at your own pace. We appreciate your contribution and are happy to wait for updates. If you have questions, feel free to ask on the Airflow Slack.


Note: This comment was drafted by an AI-assisted triage tool and may contain mistakes. Once you have addressed the points above, an Apache Airflow maintainer — a real person — will take the next look at your PR. We use this two-stage triage process so that our maintainers' limited time is spent where it matters most: the conversation with you.

Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py Outdated
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py
@potiuk
potiuk marked this pull request as draft July 28, 2026 16:11
Catching all IntegrityError was too broad: a FK violation (e.g. the
task_instance cascade-deleted between the dag_run_id lookup and the
flush) would also land here, causing the fallback UPDATE to match
zero rows while still returning 200 — silently losing the write.

Narrow the guard to unique-constraint messages per dialect
(matching the pattern in _UniqueConstraintErrorHandler) and re-raise
everything else.
@nagasrisai

nagasrisai commented Aug 4, 2026 •

Copy link
Copy Markdown
Contributor Author

@kaxil - Thanks for the pointers! I've pushed two follow-up commits:

  1. Narrowed the except IntegrityError guard to unique-constraint violations only (matching the per-dialect strings in _UniqueConstraintErrorHandler) so FK violations propagate correctly instead of being silently swallowed.
  2. Added two tests to TestXComsSetEndpoint — one that exercises the UPDATE fallback path by patching XComModel.set to raise a unique-constraint error, and one that confirms FK errors propagate as a 500.

Will also look into getting the static checks clean.

@potiuk

…m_set

Adds two tests to TestXComsSetEndpoint:

* test_xcom_set_concurrent_write_uses_update_fallback — pre-seeds the
  row, patches XComModel.set to raise a unique-constraint IntegrityError
  once, and asserts the endpoint returns 201 with the latest value
  stored via the UPDATE fallback.

* test_xcom_set_fk_integrity_error_propagates — patches XComModel.set
  to raise an FK IntegrityError and asserts it propagates as a 500
  rather than being silently swallowed.
@nagasrisai
nagasrisai marked this pull request as ready for review August 5, 2026 03:33
@nagasrisai nagasrisai changed the title fix: handle concurrent XCom writes with savepoint + UPDATE fallback to avoid 409 Handle concurrent XCom writes with savepoint and UPDATE fallback to avoid 409 Aug 8, 2026
@nagasrisai

Copy link
Copy Markdown
Contributor Author

@potiuk - Could you please mark this PR as ready for review? It's currently waiting for a reviewer. Thanks!

CC @kaxil

@nagasrisai
nagasrisai requested a review from kaxil August 17, 2026 06:47

@potiuk potiuk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for chasing this. The race is real. The task SDK client retries a POST that timed out client-side (_should_retry_api_request retries httpx.RequestError) while the first request is still running on the server. XComModel.set does DELETE then INSERT, so the second insert hits the PK, and the generic _UniqueConstraintErrorHandler turns that into a 409 the SDK never handles. Making that write succeed is a bug fix and does not need a new Execution API version.

The branch now conflicts with main, though, and the conflict changes the picture:

  1. Rebase onto #73005. task_map was removed and its length now lives in xcom.mapped_length, passed through XComModel.set(..., mapped_length=...). The savepoint here was justified by protecting the TaskMap merge, which no longer exists. As written, the UPDATE fallback (.values(value=value, dag_result=dag_result)) would also leave a stale mapped_length and timestamp on the row.

  2. Please move the fix into XComModel.set as a single upsert. airflow.utils.sqlalchemy.build_upsert_stmt already builds the dialect-specific INSERT … ON CONFLICT DO UPDATE / ON DUPLICATE KEY UPDATE, and its docstring describes this race. Replacing the DELETE+INSERT in airflow/models/xcom.py with an upsert on (dag_run_id, task_id, map_index, key), updating value, dag_result, mapped_length and timestamp:

    • fixes every caller, not only this route;
    • drops the error-string matching and the savepoint;
    • avoids a problem the savepoint approach still has on MySQL. With three or more concurrent writers (the issue reports three), the losers each hold a shared lock on the duplicate row, and their UPDATEs deadlock (1213). That rolls back the whole transaction, not just the savepoint.
  3. Don't build a 500 with e.orig in the detail. HTTPException(500, detail=f"Unexpected database error: {e.orig}") sends the raw driver message back to the caller and bypasses the [api] expose_stacktrace redaction that _DatabaseErrorHandler applies. If any IntegrityError handling remains, a bare raise lets the registered handlers produce the usual redacted response.

  4. Tests should hit the database for real. Patching XComModel.set to raise a hand-built IntegrityError with a Postgres message never enters the savepoint (or the upsert) with a real conflict, and never checks the MySQL or SQLite paths. A test that causes an actual PK conflict, for example by seeding a committed row and making the first write collide with it, would run on all three DB backends in CI. Please also move the patch / IntegrityError imports to module level.

  5. Please drop airflow-core/newsfragments/69959.bugfix.rst. As with #64614, the commit message is enough for this fix.


Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting

This branch has not been deployed

No deployments
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.

Concurrent/retried Execution API XCom writes fail with 409 due to duplicate primary key

3 participants