DG47 opened a new pull request, #74279:
URL: https://github.com/apache/airflow/pull/74279
<!-- SPDX-License-Identifier: Apache-2.0
https://www.apache.org/licenses/LICENSE-2.0 -->
`SQLExecuteQueryOperator` had no `on_kill` implementation, so killing a task
(marking it as failed in the UI, a `SIGTERM` from the executor, or an elapsed
`execution_timeout`) left the SQL statement running on the database until it
completed on its own.
This PR adds a generic, hook-level cancellation mechanism to `common.sql`
and wires it into the operator:
* `DbApiHook._run_command` keeps a reference to the cursor while
`cursor.execute()` is in flight (cleared in a `finally`), so the hook knows
which statement is currently running.
* New public method `DbApiHook.cancel_query() -> bool`. The default
implementation uses the `cancel()` extension of the DB-API cursor, falling back
to `cursor.connection.cancel()`. That covers `psycopg2`, `psycopg`, `oracledb`,
`pyodbc`, `trino`, `databricks-sql-connector` and others out of the box. When
no statement is running, or the driver has no cancel extension (e.g. `pymysql`,
`mysqlclient`), it logs a message and returns `False` instead of raising. Hooks
for databases that need a server-side command (`KILL QUERY`,
`SYSTEM$CANCEL_QUERY`, ...) can override it.
* `SQLExecuteQueryOperator.on_kill` delegates to
`self.get_db_hook().cancel_query()` and never raises, so it is safe on the kill
path (Airflow calls `on_kill` from the task process `SIGTERM` handler while
`execute()` may still be blocked inside the driver, and after
`AirflowTaskTimeout`).
Design notes, relative to the earlier attempt in #27514:
* That PR tracked query ids *after* `execute()` returned, which (as pointed
out in review) means the "running" queries had already finished. Tracking the
live cursor during `execute()` is what makes cancellation from the signal
handler actually reach a running statement.
* The state lives on the hook rather than the operator because the hook owns
the connection/cursor lifecycle; the operator only has the hook object. Hooks
with a custom `run()` that still call `self._run_command` (e.g.
`DatabricksSqlHook`) are covered automatically; hooks that bypass
`_run_command` entirely (e.g. `SnowflakeHook`) get the safe no-op until they
override `cancel_query`.
* No locks: `on_kill` runs in the main thread's signal handler, so attribute
assignment is sufficient and a lock could deadlock against the interrupted
frame.
Tests mock the hook (operator) and the DB-API cursor/connection (hook),
including the case where `cancel_query` is invoked from inside `cursor.execute`
to simulate a kill mid-statement. The `hooks/sql.pyi` public API stub gains
`cancel_query` and the operator guide documents the behaviour. No changelog
entry is added since the provider changelog is maintained by the release
manager and provider newsfragments are not used.
closes: #27314
---
* Read the **[Pull Request
Guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#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](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals))
is needed.
* When adding dependency, check compliance with the [ASF 3rd Party License
Policy](https://www.apache.org/legal/resolved.html#category-x).
* For significant user-facing changes create newsfragment:
`{pr_number}.significant.rst`, in
[airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments).
You can add this file in a follow-up commit after the PR is created so you
know the PR number.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]