Skip to content

Recover Databricks SQL statements after task retries - #72422

Closed
1fanwang wants to merge 2 commits into
apache:mainfrom
1fanwang:databricks-sql-resumability
Closed

1fanwang wants to merge 2 commits into
apache:mainfrom
1fanwang:databricks-sql-resumability

Conversation

@1fanwang

@1fanwang 1fanwang commented Sep 2, 2026 •

Copy link
Copy Markdown
Contributor

When a worker dies after submitting a Databricks SQL statement, a synchronous task retry submits the statement again. The original statement can still be running, so users can get duplicate work and cost.

Before this change, the retry submits statement-2 while statement-1 remains active. After this change, Airflow 3.3+ stores statement-1 in task state and reconnects to it. A retry also returns an already successful statement without resubmitting, and starts fresh only after failure, cancellation, closure, or a missing statement.

This uses the AIP-103 ResumableJobMixin contract only when wait_for_termination=True and deferrable=False. The statement ID reaches task state before its XCom is published, so an XCom error cannot reopen the duplicate-submission window. Fire-and-forget and deferrable execution keep their existing paths.

Testing

The driver below calls execute() twice with the real task-state accessor and supervisor messages. Its local Databricks hook forces statement-ID XCom publication to fail after the first submission.

Live task-state and XCom ordering proof
from __future__ import annotations

import importlib
import importlib.util
import os
import sys
from pathlib import Path
from typing import Any
from unittest import mock
from uuid import UUID

from airflow.providers.databricks.hooks.databricks import SQLStatementState
from airflow.sdk._shared.state import TaskScope
from airflow.sdk.execution_time import task_runner
from airflow.sdk.execution_time.comms import GetTaskStateStore, SetTaskStateStore, TaskStateStoreResult
from airflow.sdk.execution_time.context import TaskStateStoreAccessor

MODULE_NAME = "airflow.providers.databricks.operators.databricks"
if source_path := os.environ.get("DATABRICKS_OPERATOR_SOURCE"):
    module_spec = importlib.util.spec_from_file_location(MODULE_NAME, Path(source_path))
    if module_spec is None or module_spec.loader is None:
        raise RuntimeError(f"cannot load {source_path}")
    databricks_module = importlib.util.module_from_spec(module_spec)
    sys.modules[MODULE_NAME] = databricks_module
    module_spec.loader.exec_module(databricks_module)
else:
    databricks_module = importlib.import_module(MODULE_NAME)

DatabricksSQLStatementsOperator = databricks_module.DatabricksSQLStatementsOperator


class LocalDatabricksHook:
    def __init__(self) -> None:
        self.submissions = 0
        self.status_checks = 0

    def post_sql_statement(self, json: dict[str, Any]) -> str:
        self.submissions += 1
        return f"statement-{self.submissions}"

    def get_sql_statement_state(self, statement_id: str) -> SQLStatementState:
        self.status_checks += 1
        if self.status_checks == 1:
            return SQLStatementState("RUNNING")
        return SQLStatementState("SUCCEEDED")


def main(expected_submissions: int) -> None:
    stored: dict[str, Any] = {}

    def send(message: Any) -> TaskStateStoreResult | None:
        if isinstance(message, GetTaskStateStore):
            value = stored.get(message.key)
            return TaskStateStoreResult(value=value) if value is not None else None
        if isinstance(message, SetTaskStateStore):
            stored[message.key] = message.value
        return None

    task_state_store = TaskStateStoreAccessor(
        ti_id=UUID("00000000-0000-0000-0000-000000000001"),
        scope=TaskScope(dag_id="proof", run_id="run", task_id="sql"),
    )
    first_ti = mock.MagicMock(stats_tags={})
    first_ti.xcom_push.side_effect = RuntimeError("xcom unavailable")
    first_context = {"task_state_store": task_state_store, "ti": first_ti}
    retry_ti = mock.MagicMock(stats_tags={})
    retry_context = {"task_state_store": task_state_store, "ti": retry_ti}
    hook = LocalDatabricksHook()

    with (
        mock.patch.object(task_runner, "SUPERVISOR_COMMS", mock.Mock(send=mock.Mock(side_effect=send)), create=True),
        mock.patch.object(DatabricksSQLStatementsOperator, "_get_hook", return_value=hook),
    ):
        first = DatabricksSQLStatementsOperator(
            task_id="sql",
            statement="SELECT 1",
            warehouse_id="warehouse",
            polling_period_seconds=0,
        )
        try:
            first.execute(first_context)
        except RuntimeError as error:
            print(f"first_attempt={error}")
        print(f"stored_after_first={stored}")

        retry = DatabricksSQLStatementsOperator(
            task_id="sql",
            statement="SELECT 1",
            warehouse_id="warehouse",
            polling_period_seconds=0,
        )
        retry.execute(retry_context)

    print(f"stored_after_retry={stored}")
    print(f"submissions={hook.submissions}")
    print(f"retry_statement_id={retry.statement_id}")
    print(f"retry_xcom_calls={retry_ti.xcom_push.call_args_list}")
    if hook.submissions != expected_submissions:
        raise RuntimeError(f"expected {expected_submissions} submission")


if __name__ == "__main__":
    main(expected_submissions=int(sys.argv[1]))

Save the driver as dev/databricks_xcom_ordering_e2e.py, then run it against the previous PR head:

mkdir -p dev/databricks-xcom-ordering-pre-fix
git archive 82cddda96f6265acad56a285a54ad2f1a05df177 \
  providers/databricks/src/airflow/providers/databricks/operators/databricks.py \
  | tar -x -C dev/databricks-xcom-ordering-pre-fix
DATABRICKS_OPERATOR_SOURCE="$PWD/dev/databricks-xcom-ordering-pre-fix/providers/databricks/src/airflow/providers/databricks/operators/databricks.py" \
AIRFLOW_HOME="$PWD/dev/databricks-xcom-ordering-red-airflow-home" \
AIRFLOW__CORE__LOAD_EXAMPLES=False \
.venv/bin/uv run --project providers/databricks python dev/databricks_xcom_ordering_e2e.py 1
first_attempt=xcom unavailable
stored_after_first={}
stored_after_retry={'databricks_sql_statement_id': 'statement-2'}
submissions=2
retry_statement_id=statement-2
retry_xcom_calls=[call(key='statement_id', value='statement-2')]
RuntimeError: expected 1 submission

Run the same driver against this branch:

AIRFLOW_HOME="$PWD/dev/databricks-xcom-ordering-green-airflow-home" \
AIRFLOW__CORE__LOAD_EXAMPLES=False \
.venv/bin/uv run --project providers/databricks python dev/databricks_xcom_ordering_e2e.py 1
first_attempt=xcom unavailable
stored_after_first={'databricks_sql_statement_id': 'statement-1'}
Reconnecting to existing job external_id=statement-1 external_id_key=databricks_sql_statement_id status=RUNNING
stored_after_retry={'databricks_sql_statement_id': 'statement-1'}
submissions=1
retry_statement_id=statement-1
retry_xcom_calls=[call(key='statement_id', value='statement-1')]

Please check the type of change your PR introduces:

  • Bugfix
  • Feature
  • New Provider
  • Improvement
  • Documentation Fix
  • Refactoring
  • Other

Relevant Issue(s)

None.

Description

See above.

How did you test it?

See the live proof above.

Did you add documentation?

Yes. The SQL statements guide covers the recovery behavior and version floor.

Does this introduce a breaking change?

No.

Checklist

  • I have checked that there are no existing pull requests for the same change.
  • I have added a commit message that describes my changes.
  • I have added tests that prove my fix is effective or that my feature works.
  • I have updated the documentation accordingly.

Do you use Generative AI to contribute to this PR?

  • No
  • Yes, I used GitHub Copilot CLI (GPT-5.6 Sol) to assist with this pull request.

AI-assisted review checklist

  • I have manually reviewed all AI-generated code and verified its correctness.
  • I have added comprehensive tests for AI-generated code paths.
  • I have verified that no sensitive data or secrets were included in AI prompts.
  • I understand that I am responsible for the entire content of this PR.

PR Checklist

  • I have locally rebased my branch onto the latest main.
  • I have run relevant tests and they pass.
  • I have run pre-commit checks and they pass.
  • I have updated documentation where applicable.

Synchronous SQL statements could be submitted twice when a worker crashed after submission because retries had no persisted external ID.

Generated-by: GitHub Copilot CLI (GPT-5.6 Sol)
Signed-off-by: 1fanwang <1fannnw@gmail.com>
Persist the Databricks statement ID before publishing its XCom so a worker failure can reconnect on retry.

Generated-by: GitHub Copilot CLI (GPT-5.6 Sol)
Signed-off-by: 1fanwang <1fannnw@gmail.com>
This was referenced Sep 25, 2026
@potiuk potiuk added the closed because of open PR limit Closed as a one-time step of introducing the open pull request limit label Sep 25, 2026
@potiuk

potiuk commented Sep 25, 2026

Copy link
Copy Markdown
Member

Hello @1fanwang - 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 34 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 gh pr reopen <PR_NUMBER> --repo apache/airflow. Reopen the ones you are ready to follow through - keep them rebased, respond to review comments and fix failing checks.

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

@potiuk potiuk closed this Sep 25, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providers closed because of open PR limit Closed as a one-time step of introducing the open pull request limit kind:documentation provider:databricks

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants