mirror of
https://github.com/apache/superset.git
synced 2026-09-01 13:01:33 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e878fe7d40 |
@@ -43,6 +43,7 @@ from superset.db_engine_specs.base import (
|
|||||||
from superset.db_engine_specs.hive import HiveEngineSpec
|
from superset.db_engine_specs.hive import HiveEngineSpec
|
||||||
from superset.errors import ErrorLevel, SupersetError, SupersetErrorType
|
from superset.errors import ErrorLevel, SupersetError, SupersetErrorType
|
||||||
from superset.exceptions import OAuth2Error
|
from superset.exceptions import OAuth2Error
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
from superset.utils import json
|
from superset.utils import json
|
||||||
from superset.utils.core import get_user_agent, QuerySource
|
from superset.utils.core import get_user_agent, QuerySource
|
||||||
from superset.utils.network import is_hostname_valid, is_port_open
|
from superset.utils.network import is_hostname_valid, is_port_open
|
||||||
@@ -775,6 +776,11 @@ class DatabricksPythonConnectorEngineSpec(DatabricksDynamicBaseEngineSpec):
|
|||||||
|
|
||||||
parameters_schema = DatabricksPythonConnectorSchema()
|
parameters_schema = DatabricksPythonConnectorSchema()
|
||||||
|
|
||||||
|
# The driver only populates `cursor.active_command_id` once a statement has
|
||||||
|
# actually been executed, so the cancel id can't be captured up front like
|
||||||
|
# it can for engines where it's tied to the session rather than the query.
|
||||||
|
has_query_id_before_execute = False
|
||||||
|
|
||||||
sqlalchemy_uri_placeholder = (
|
sqlalchemy_uri_placeholder = (
|
||||||
"databricks://token:{access_token}@{host}:{port}?http_path={http_path}"
|
"databricks://token:{access_token}@{host}:{port}?http_path={http_path}"
|
||||||
"&catalog={default_catalog}&schema={default_schema}"
|
"&catalog={default_catalog}&schema={default_schema}"
|
||||||
@@ -957,6 +963,66 @@ class DatabricksPythonConnectorEngineSpec(DatabricksDynamicBaseEngineSpec):
|
|||||||
|
|
||||||
return uri, connect_args
|
return uri, connect_args
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def get_cancel_query_id(cls, cursor: Any, query: Query) -> str | None:
|
||||||
|
"""
|
||||||
|
Capture a cancel id for the query that was just executed on ``cursor``.
|
||||||
|
|
||||||
|
Cancellation is only supported when the connection uses Databricks'
|
||||||
|
Statement Execution API (SEA) backend, i.e. ``use_sea=True`` was set in
|
||||||
|
the database's connection parameters. A SEA statement id is a plain
|
||||||
|
string that can be cancelled later from a brand-new cursor/connection,
|
||||||
|
which matches how Superset issues cancellation (a fresh connection, not
|
||||||
|
the one that ran the query).
|
||||||
|
|
||||||
|
The default (Thrift) backend has no equivalent public mechanism: the
|
||||||
|
driver can only cancel a Thrift command via the live ``CommandId`` held
|
||||||
|
by the executing cursor, which includes a secret that's never exposed
|
||||||
|
through any public/documented accessor (``cursor.query_id`` only
|
||||||
|
returns the operation's GUID). Rather than reconstructing that secret
|
||||||
|
from the driver's private internals, this returns ``None`` for
|
||||||
|
Thrift-backed connections so the "stop query" action fails explicitly
|
||||||
|
instead of silently no-oping.
|
||||||
|
|
||||||
|
:param cursor: Cursor instance in which the query was just executed
|
||||||
|
:param query: Query instance
|
||||||
|
:return: SEA statement id, or None if cancellation isn't supported
|
||||||
|
"""
|
||||||
|
command_id = getattr(cursor, "active_command_id", None)
|
||||||
|
if command_id is None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
session = getattr(getattr(cursor, "connection", None), "session", None)
|
||||||
|
if not getattr(session, "use_sea", False):
|
||||||
|
return None
|
||||||
|
|
||||||
|
return command_id.to_sea_statement_id()
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def cancel_query(cls, cursor: Any, query: Query, cancel_query_id: str) -> bool:
|
||||||
|
"""
|
||||||
|
Cancel a query in the underlying database.
|
||||||
|
|
||||||
|
Only implemented for SEA (Statement Execution API) connections; see
|
||||||
|
``get_cancel_query_id``. Any error raised while attempting the cancel
|
||||||
|
(e.g. the fresh cursor/connection itself failing) is allowed to
|
||||||
|
propagate rather than being reported as a successful or failed cancel.
|
||||||
|
|
||||||
|
:param cursor: New cursor instance to the db of the query
|
||||||
|
:param query: Query instance
|
||||||
|
:param cancel_query_id: SEA statement id, as returned by
|
||||||
|
``get_cancel_query_id``
|
||||||
|
:return: True if the cancel request was submitted successfully
|
||||||
|
"""
|
||||||
|
if not cls.validate_cancel_query_id(cancel_query_id, r"^[a-zA-Z0-9-]+$"):
|
||||||
|
return False
|
||||||
|
|
||||||
|
from databricks.sql.backend.types import CommandId
|
||||||
|
|
||||||
|
cursor.active_command_id = CommandId.from_sea_statement_id(cancel_query_id)
|
||||||
|
cursor.cancel()
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
# TODO: remove once we've upgraded to SQLAlchemy>=2.0 and databricks-sql-python>=3.x
|
# TODO: remove once we've upgraded to SQLAlchemy>=2.0 and databricks-sql-python>=3.x
|
||||||
monkeypatch_dialect()
|
monkeypatch_dialect()
|
||||||
|
|||||||
@@ -1092,3 +1092,124 @@ def test_get_oauth2_fresh_token_python(
|
|||||||
},
|
},
|
||||||
timeout=30.0,
|
timeout=30.0,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_has_query_id_before_execute_is_false() -> None:
|
||||||
|
"""
|
||||||
|
The cancel id can only be read off the cursor after a statement has been
|
||||||
|
executed (``cursor.active_command_id`` is unset before that), so the id
|
||||||
|
must be captured post-execute rather than up front.
|
||||||
|
"""
|
||||||
|
assert DatabricksPythonConnectorEngineSpec.has_query_id_before_execute is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_get_cancel_query_id_sea(mocker: MockerFixture) -> None:
|
||||||
|
"""
|
||||||
|
On a SEA (Statement Execution API) connection, the cancel id is the
|
||||||
|
statement id exposed by the driver's ``CommandId``.
|
||||||
|
"""
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
|
|
||||||
|
cursor = mocker.MagicMock()
|
||||||
|
cursor.active_command_id.to_sea_statement_id.return_value = "01ecc35f-abcd"
|
||||||
|
cursor.connection.session.use_sea = True
|
||||||
|
|
||||||
|
query = Query()
|
||||||
|
assert (
|
||||||
|
DatabricksPythonConnectorEngineSpec.get_cancel_query_id(cursor, query)
|
||||||
|
== "01ecc35f-abcd"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_get_cancel_query_id_no_active_command(mocker: MockerFixture) -> None:
|
||||||
|
"""
|
||||||
|
Before a statement has been executed on the cursor, there's nothing to
|
||||||
|
build a cancel id from.
|
||||||
|
"""
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
|
|
||||||
|
cursor = mocker.MagicMock()
|
||||||
|
cursor.active_command_id = None
|
||||||
|
cursor.connection.session.use_sea = True
|
||||||
|
|
||||||
|
query = Query()
|
||||||
|
assert (
|
||||||
|
DatabricksPythonConnectorEngineSpec.get_cancel_query_id(cursor, query) is None
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_get_cancel_query_id_thrift_backend_returns_none(
|
||||||
|
mocker: MockerFixture,
|
||||||
|
) -> None:
|
||||||
|
"""
|
||||||
|
When the connection was not established with ``use_sea=True`` (i.e. it's
|
||||||
|
using the default Thrift backend), no cancel id is captured -- cancelling
|
||||||
|
a Thrift command requires a secret that's never exposed through any
|
||||||
|
public driver API, so we deliberately don't attempt it.
|
||||||
|
"""
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
|
|
||||||
|
cursor = mocker.MagicMock()
|
||||||
|
cursor.active_command_id.to_sea_statement_id.return_value = "01ecc35f-abcd"
|
||||||
|
cursor.connection.session.use_sea = False
|
||||||
|
|
||||||
|
query = Query()
|
||||||
|
assert (
|
||||||
|
DatabricksPythonConnectorEngineSpec.get_cancel_query_id(cursor, query) is None
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_cancel_query_sea_success(mocker: MockerFixture) -> None:
|
||||||
|
"""
|
||||||
|
``cancel_query`` sets the reconstructed SEA ``CommandId`` on the fresh
|
||||||
|
cursor and delegates to the driver's own ``cursor.cancel()``.
|
||||||
|
"""
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
|
|
||||||
|
command_id = mocker.patch(
|
||||||
|
"databricks.sql.backend.types.CommandId.from_sea_statement_id"
|
||||||
|
)
|
||||||
|
cursor = mocker.MagicMock()
|
||||||
|
|
||||||
|
query = Query()
|
||||||
|
assert (
|
||||||
|
DatabricksPythonConnectorEngineSpec.cancel_query(cursor, query, "01ecc35f-abcd")
|
||||||
|
is True
|
||||||
|
)
|
||||||
|
command_id.assert_called_once_with("01ecc35f-abcd")
|
||||||
|
assert cursor.active_command_id == command_id.return_value
|
||||||
|
cursor.cancel.assert_called_once_with()
|
||||||
|
|
||||||
|
|
||||||
|
def test_cancel_query_invalid_id_returns_false(mocker: MockerFixture) -> None:
|
||||||
|
"""
|
||||||
|
A malformed cancel id is rejected before touching the cursor at all.
|
||||||
|
"""
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
|
|
||||||
|
cursor = mocker.MagicMock()
|
||||||
|
query = Query()
|
||||||
|
assert (
|
||||||
|
DatabricksPythonConnectorEngineSpec.cancel_query(
|
||||||
|
cursor, query, "'; DROP TABLE foo; --"
|
||||||
|
)
|
||||||
|
is False
|
||||||
|
)
|
||||||
|
cursor.cancel.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
def test_cancel_query_propagates_errors(mocker: MockerFixture) -> None:
|
||||||
|
"""
|
||||||
|
If the cancel attempt itself fails (e.g. the fresh cursor/connection
|
||||||
|
errors), the error must surface rather than being reported as a failed-
|
||||||
|
but-handled cancel.
|
||||||
|
"""
|
||||||
|
from superset.models.sql_lab import Query
|
||||||
|
|
||||||
|
mocker.patch("databricks.sql.backend.types.CommandId.from_sea_statement_id")
|
||||||
|
cursor = mocker.MagicMock()
|
||||||
|
cursor.cancel.side_effect = RuntimeError("connection reset")
|
||||||
|
|
||||||
|
query = Query()
|
||||||
|
with pytest.raises(RuntimeError):
|
||||||
|
DatabricksPythonConnectorEngineSpec.cancel_query(cursor, query, "01ecc35f-abcd")
|
||||||
|
|||||||
Reference in New Issue
Block a user