Skip to content

Add durable execution to RedshiftDataOperator#69530

Merged
amoghrajesh merged 4 commits into
apache:mainfrom
astronomer:redshift-crash-recovery
Jul 13, 2026
Merged

Add durable execution to RedshiftDataOperator#69530
amoghrajesh merged 4 commits into
apache:mainfrom
astronomer:redshift-crash-recovery

Conversation

@amoghrajesh

@amoghrajesh amoghrajesh commented Jul 7, 2026

Copy link
Copy Markdown
Contributor

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

What

RedshiftDataOperator submits a statement to the Redshift Data API and polls it to completion on the worker. On a worker crash or preemption mid-poll, Airflow retries the task by calling execute() again which resubmits the SQL from scratch, since nothing about the in-flight statement id is persisted across attempts.

Current behaviour

A retry after a crash always resubmits the full SQL, even if the original statement is still running (or already finished) in Redshift.

  • For non-idempotent SQL (INSERT, COPY, UPDATE, CREATE TABLE) this risks duplicate writes;
  • For expensive queries it's wasted compute, since the orphaned original execution keeps running with nobody polling it.

This also applies, in a narrower window, to wait_for_completion=False: even though that mode does not poll at all, a retry after a successful submission still resubmits today, since the id is never persisted regardless of whether the task waits.

Proposed change

Adds ResumableJobMixin support (Airflow 3.3+) to RedshiftDataOperator, following the same pattern already done for the Databricks operators and SnowflakeSqlApiOperator. Before polling begins, the submitted statement id is persisted to task_state_store. On retry, the operator reads it back and:

  • reconnects and keeps polling if the statement is still running
  • returns immediately without resubmitting if it already succeeded
  • submits the SQL fresh if the statement failed, or its id has expired and is no longer found (ClientError)

durable=True is the default; set durable=False to keep the old "always submit fresh on retry" behavior. deferrable=True takes precedence over durable, the Triggerer already tracks the statement across the wait in that mode.

Critically, the persistence-on-submit design means the wait_for_completion=False case is protected too, not just the blocking-wait case & the id is written to task_state_store immediately after submission regardless of whether the task polls for it.

Changes of Note

  • submit_job now always calls hook.execute_query(..., wait_for_completion=False, ...) and hence the mixin, not the hook, owns polling. The statement id is set on the operator immediately after submission (not after the wait completes), which also fixes a pre-existing on_kill gap: today, killing the task mid-wait doesn't cancel the statement because self.statement_id isn't set yet.
  • get_job_status calls describe_statement directly rather than reusing check_query_is_finished/parse_statement_response, since those raise on failure states rather than returning a status string, which would break the mixin's status-in/decision-out contract.
  • is_job_active uses a negative test (not in (FINISHED, *FAILURE_STATES, NOT_FOUND)) rather than a positive running-states allowlist, deliberately avoiding a real pitfall found in RedshiftDataTrigger.is_still_running, which uses a positive allowlist and would misclassify any future unlisted Redshift status as "done."
  • get_sql_results (the actual return value fetch) was moved into poll_until_complete rather than left solely in get_job_result, because the mixin calls only poll_until_complete on reconnect - get_job_result is skipped entirely on that path. A first pass missed returning the fetched value from poll_until_complete itself, which silently turned execute()'s return value into None on reconnect; caught via the durable test suite and fixed.
  • Verified there was no missing-state prerequisite bugfix needed here (unlike the Databricks port, which needed BLOCKED/WAITING_FOR_RETRY added to a state allowlist first) - Redshift's RUNNING_STATES/FINISHED_STATE/FAILURE_STATES are already complete against the real AWS API's Status enum.

User implications / backcompat

No breaking change. durable defaults to True on Airflow 3.3+; on earlier versions it's a no-op stub and the operator always submits fresh, exactly as before. If task_state_store isn't available at runtime, the operator logs that crash recovery is disabled and falls back to the same fresh-submit behavior.

Testing

  • Create an aws_default connection
  • Create a redshift cluster - I created one with tickit data loaded and am trying to run some queries using this DAG:
from datetime import datetime, timedelta

from airflow.providers.amazon.aws.operators.redshift_data import RedshiftDataOperator
from airflow.sdk import DAG

with DAG(
    dag_id="try_redshift_data",
    schedule=None,
    start_date=datetime(2024, 1, 1),
    catchup=False,
) as dag:
    query_tickit = RedshiftDataOperator(
        task_id="query_tickit",
        aws_conn_id="aws_default",
        cluster_identifier="redshift-cluster-demo",
        database="dev",
        db_user="awsuser",
        sql="SELECT COUNT(*) FROM sales a, sales b, listing c WHERE a.salesid < 20000 AND b.salesid < 20000;",
        wait_for_completion=True,
        deferrable=False,
        retries=3,
        retry_delay=timedelta(seconds=5),
    )

Before my changes

Try 1:

image image

Worker was killed above, once worker comes back up, and a new redshift query gets submitted:

image image

Waste / Duplicate submission.

After my changes

Try 1:

image image

But now its saved to task state store:

image

Since I am using a custom backend, this is the data:

[Breeze:3.10.20] root@7ce59ff19824:/opt/airflow$ cat /tmp/airflow_state/ti_query_tickit/redshift_statement_id.json
"dd5d4af5-c05d-401f-9958-2d0622a8764a"[Breeze:3.10.20] root@7ce59ff19824:/opt/airflow

Try 2 when worker comes up:

image
  • 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.

@amoghrajesh

Copy link
Copy Markdown
Contributor Author

Added a dont-merge label until I can get my hands on testing this on a real cluster

@amoghrajesh
amoghrajesh requested a review from uranusjr July 7, 2026 07:15
@amoghrajesh
amoghrajesh marked this pull request as draft July 7, 2026 07:15
@amoghrajesh amoghrajesh moved this from Backlog to In progress in Durable / Crash-Safe Execution Jul 8, 2026
@amoghrajesh

Copy link
Copy Markdown
Contributor Author

@o-nikolas I was able to test it, removed dont-merge label and marked it ready for review. Check PR desc for testing details.

@amoghrajesh
amoghrajesh marked this pull request as ready for review July 8, 2026 12:18
@amoghrajesh

Copy link
Copy Markdown
Contributor Author

@vincbeck / @o-nikolas this one's ready now, will love to get some reviews here if you have some time.

@amoghrajesh
amoghrajesh requested a review from vincbeck July 10, 2026 05:58
@amoghrajesh

Copy link
Copy Markdown
Contributor Author

Thanks for the review @vincbeck! Unrelated failures, already handled by: #69798. Merging this one.

@amoghrajesh
amoghrajesh merged commit 5d1e204 into apache:main Jul 13, 2026
76 of 81 checks passed
@amoghrajesh
amoghrajesh deleted the redshift-crash-recovery branch July 13, 2026 07:00
@github-project-automation github-project-automation Bot moved this from In progress to Done in Durable / Crash-Safe Execution Jul 13, 2026
joshuabvarghese pushed a commit to joshuabvarghese/airflow that referenced this pull request Jul 16, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

Development

Successfully merging this pull request may close these issues.

2 participants