Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
# Release History

# Unreleased
- Add the default-enabled, kernel-only `enable_geospatial_support` connection option. GEOMETRY / GEOGRAPHY results are exposed as `{"srid": int, "wkb": bytes}` values when true or WKT / EWKT strings when false.

# 4.6.0 (2026-09-24)
- Upgrade Databricks SQL Kernel to 1.1.0; the kernel dependency is now stable and no longer experimental.
- Transparently auto-recover Thrift connections to Reyden / Real-Time warehouses: when a warehouse rejects the default Thrift protocol (SQLSTATE `KP001`), the session is re-opened on the kernel backend and the warehouse is remembered so later connections skip Thrift. Applies only when no backend was chosen explicitly.
Expand Down
1 change: 1 addition & 0 deletions CONNECTION_PARAMETERS.md
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ to change without notice.
| `max_download_threads` | `int` | ✅ | ❌ | `10` | Worker threads for cloud-fetch downloads. Not forwarded to the kernel. |
| `enable_query_result_lz4_compression` | `bool` | ✅ | ❌ | `True` | LZ4-compress result payloads. Not forwarded; the kernel handles compression internally. |
| `_disable_pandas` | `bool` | ✅ | ✅ | `False` | Skip the pandas-based Arrow→row deserialization and materialize rows directly with PyArrow. This is a **Python-side** result-conversion toggle, not a wire option: the kernel returns results as Arrow (`RecordBatch`es) and the connector runs the *same* `_convert_arrow_table` for both backends, so the flag is honored on the kernel path too. Affects only row fetches (`fetchone`/`fetchmany`/`fetchall`); the `fetch*_arrow` methods return the Arrow table unchanged regardless of this flag. |
| `enable_geospatial_support` | `bool` | ❌ | ✅ | `True` | Return GEOMETRY / GEOGRAPHY as `{"srid": int, "wkb": bytes}` values when `True`, or as WKT / EWKT strings when `False`. This is a local result conversion and is never forwarded to SEA. |
| `_use_arrow_native_complex_types` | `bool` | ✅ | ✅ | `True` | Return `ARRAY`/`MAP`/`STRUCT` as native Arrow types instead of JSON strings. Forwarded to the kernel. |
| `_use_arrow_native_decimals` | `bool` | ✅ | ❌ | `True` | Thrift wire encoding for `DECIMAL`: `True` → native Arrow `decimal128`, `False` → Arrow string. **No value-level effect**, though: the connector unconditionally re-casts the column back to `decimal128` (`convert_decimals_in_arrow_table`, `thrift_backend.py`), so both `fetchall()` and `fetchall_arrow()` yield `Decimal` / `decimal128(p,s)` either way (verified live). Not forwarded to the kernel, which always returns native Arrow decimals. |
| `_use_arrow_native_timestamps` | `bool` | ✅ | ❌ | `True` | Thrift wire encoding for `TIMESTAMP`: `True` → native Arrow timestamp (→ Python `datetime`), `False` → Arrow string (→ Python **`str`**). **Unlike decimals there is no re-cast**, so `False` genuinely surfaces strings — and `cursor.description` still reports the type code as `'timestamp'`, a mismatch to watch for (verified live). Note the connector always also sends the `spark.thriftserver.arrowBasedRowSet.timestampAsString=false` conf, but the `timestampAsArrow=False` flag wins. Not forwarded to the kernel, which always returns native Arrow timestamps. |
Expand Down
2 changes: 1 addition & 1 deletion KERNEL_REV
Original file line number Diff line number Diff line change
@@ -1 +1 @@
80f2aee7d884994d7b0af9a9ea6078872859a9cd
b7e9310b27be16a6c42490e58c4b5c4a525320bc
30 changes: 30 additions & 0 deletions src/databricks/sql/backend/kernel/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,26 @@ def _kernel_session_accepts_kwarg(name: str) -> bool:
return name in params


def _kernel_geospatial_kwargs(value: bool) -> Dict[str, bool]:
"""Build the default-enabled native geospatial-support kwarg.

The value is always passed explicitly so the driver and kernel contracts
cannot drift. Older kernel wheels do not declare
``enable_geospatial_support`` and must fail clearly rather than silently
return a different public value shape.
"""
if not isinstance(value, bool):
raise ValueError(
"enable_geospatial_support must be a bool; " f"got {type(value).__name__}"
)
if not _kernel_session_accepts_kwarg("enable_geospatial_support"):
raise NotSupportedError(
"enable_geospatial_support requires a newer databricks-sql-kernel "
"wheel that exposes geospatial result representation support."
)
return {"enable_geospatial_support": value}


def _kernel_telemetry_kwargs(options: Dict[str, Any]) -> Dict[str, Any]:
"""Build phase-7 telemetry/system kwargs for ``databricks_sql_kernel.Session``.

Expand Down Expand Up @@ -262,6 +282,12 @@ def __init__(
# The kernel binding owns type and range validation.
self._request_timeout_secs = kwargs.get("request_timeout_secs")
self._max_connections = kwargs.get("max_connections")
# Default-enabled native GEOMETRY / GEOGRAPHY support. True requests
# the canonical Arrow struct and surfaces as ``{"srid": int, "wkb":
# bytes}`` through pyarrow; False requests WKT / EWKT strings. This is
# intentionally separate from ``session_configuration``: it is never
# forwarded to SEA.
self._enable_geospatial_support = kwargs.get("enable_geospatial_support", True)
# Kernel telemetry phase 7 adds binding/runtime identity and
# telemetry config kwargs directly to ``databricks_sql_kernel.Session``.
self._telemetry_options = kwargs.get("telemetry_options") or {}
Expand Down Expand Up @@ -379,6 +405,9 @@ def open_session(
# kernel's ``retry_*`` kwargs. Empty when at defaults.
retry_kwargs = _kernel_retry_kwargs(self._retry_options)
telemetry_kwargs = _kernel_telemetry_kwargs(self._telemetry_options)
geospatial_kwargs = _kernel_geospatial_kwargs(
self._enable_geospatial_support
)
max_connections_kwargs: Dict[str, Any] = {}
if _kernel_session_accepts_kwarg("max_connections"):
max_connections_kwargs["max_connections"] = self._max_connections
Expand Down Expand Up @@ -426,6 +455,7 @@ def open_session(
**tls_kwargs,
**retry_kwargs,
**telemetry_kwargs,
**geospatial_kwargs,
**max_connections_kwargs,
**http_headers_kwargs,
)
Expand Down
7 changes: 7 additions & 0 deletions src/databricks/sql/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,13 @@ def __init__(
decision. This is an intentional divergence from the
Thrift/SEA paths, where an explicit ``True`` can still
be suppressed by the feature flag.
:param enable_geospatial_support: `bool`, optional (default is True)
Kernel backend only. Controls the public representation of
``GEOMETRY`` and ``GEOGRAPHY`` result values. ``True`` returns
``{"srid": int, "wkb": bytes}``; ``False`` returns WKT / EWKT
strings (for example ``"SRID=4326;POINT(1 2)"``). The
conversion is local to the kernel/driver and this option is
never sent to the SQL Execution API.
:param use_hybrid_disposition: `bool`, optional (default is False)
Use the hybrid disposition instead of the inline disposition.
:param server_hostname: Databricks instance host name.
Expand Down
1 change: 1 addition & 0 deletions src/databricks/sql/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,7 @@ def _create_backend(
retry_options=kernel_retry_options,
request_timeout_secs=kwargs.get("_socket_timeout"),
max_connections=kwargs.get("_pool_maxsize") or None,
enable_geospatial_support=kwargs.get("enable_geospatial_support", True),
telemetry_options=kernel_telemetry_options,
)

Expand Down
93 changes: 92 additions & 1 deletion tests/unit/test_kernel_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -392,6 +392,94 @@ def fake_session(**kw):
assert captured["max_connections"] == max_connections


@pytest.mark.parametrize("enabled", [True, False])
def test_open_session_passes_geospatial_representation_to_kernel(
monkeypatch, enabled
):
captured = {}

def fake_session(*, enable_geospatial_support=True, **kw):
captured["enable_geospatial_support"] = enable_geospatial_support
sess = MagicMock()
sess.session_id = "sess-id"
return sess

monkeypatch.setattr(kernel_client._kernel, "Session", fake_session)
c = kernel_client.KernelDatabricksClient(
server_hostname="example.cloud.databricks.com",
http_path="/sql/1.0/warehouses/abc",
auth_provider=AccessTokenAuthProvider("dapi-test"),
ssl_options=None,
enable_geospatial_support=enabled,
)

c.open_session(session_configuration=None, catalog=None, schema=None)

assert captured["enable_geospatial_support"] is enabled


def test_open_session_enables_geospatial_support_by_default(monkeypatch):
captured = {}

def fake_session(**kw):
captured.update(kw)
sess = MagicMock()
sess.session_id = "sess-id"
return sess

monkeypatch.setattr(kernel_client._kernel, "Session", fake_session)
c = kernel_client.KernelDatabricksClient(
server_hostname="example.cloud.databricks.com",
http_path="/sql/1.0/warehouses/abc",
auth_provider=AccessTokenAuthProvider("dapi-test"),
ssl_options=None,
)

c.open_session(session_configuration=None, catalog=None, schema=None)

assert captured["enable_geospatial_support"] is True


def test_open_session_rejects_explicit_geospatial_option_with_old_kernel(
monkeypatch,
):
def fake_session_without_geospatial(
host,
http_path,
*,
catalog=None,
schema=None,
session_conf=None,
complex_types_as_json=False,
intervals_as_string=False,
request_timeout_secs=None,
auth_type=None,
access_token=None,
):
sess = MagicMock()
sess.session_id = "sess-id"
return sess

monkeypatch.setattr(
kernel_client._kernel, "Session", fake_session_without_geospatial
)
c = kernel_client.KernelDatabricksClient(
server_hostname="example.cloud.databricks.com",
http_path="/sql/1.0/warehouses/abc",
auth_provider=AccessTokenAuthProvider("dapi-test"),
ssl_options=None,
enable_geospatial_support=False,
)

with pytest.raises(NotSupportedError, match="newer databricks-sql-kernel"):
c.open_session(session_configuration=None, catalog=None, schema=None)


def test_geospatial_option_rejects_non_bool():
with pytest.raises(ValueError, match="must be a bool"):
kernel_client._kernel_geospatial_kwargs("false")


def test_open_session_passes_phase_7_telemetry_kwargs_to_kernel(monkeypatch):
"""Kernel telemetry phase 7 added binding/runtime identity and
telemetry config kwargs to ``databricks_sql_kernel.Session``."""
Expand Down Expand Up @@ -490,6 +578,7 @@ def fake_session_without_optional_kwargs(
catalog=None,
schema=None,
session_conf=None,
enable_geospatial_support=True,
complex_types_as_json=False,
intervals_as_string=False,
request_timeout_secs=None,
Expand Down Expand Up @@ -589,7 +678,9 @@ def raise_value_error(_obj):
kwargs = kernel_client._kernel_telemetry_kwargs(
{"enable_telemetry": True, "telemetry_batch_size": 17}
)
assert kwargs == {}, f"expected no phase-7 kwargs when signature unreadable, got {kwargs}"
assert (
kwargs == {}
), f"expected no phase-7 kwargs when signature unreadable, got {kwargs}"


def test_execute_command_forwards_parameters_to_bind_param():
Expand Down
40 changes: 39 additions & 1 deletion tests/unit/test_kernel_result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,11 +43,12 @@ def close(self):
self.closed = True


def _make_rs(handle) -> KernelResultSet:
def _make_rs(handle, *, disable_pandas=False) -> KernelResultSet:
# The base ResultSet __init__ takes a `connection` ref it never
# actually dereferences during these buffer tests, so a Mock is
# fine.
connection = MagicMock()
connection.disable_pandas = disable_pandas
backend = MagicMock()
return KernelResultSet(
connection=connection,
Expand Down Expand Up @@ -94,6 +95,43 @@ def test_fetchall_arrow_drains_all_batches(int_schema):
assert rs.has_more_rows is False


def test_geospatial_string_and_binary_values_keep_logical_type():
wkb = bytes.fromhex("0101000000000000000000F03F0000000000000040")
geo_metadata = {
b"databricks.type_name": b"GEOMETRY",
b"databricks.type_text": b"GEOMETRY(ANY)",
}

string_schema = pa.schema([pa.field("g", pa.string(), metadata=geo_metadata)])
string_batch = pa.RecordBatch.from_arrays(
[pa.array(["SRID=4326;POINT(1 2)", None], type=pa.string())],
schema=string_schema,
)
string_rows = _make_rs(
_FakeKernelHandle(string_schema, [string_batch]), disable_pandas=True
).fetchall()
assert [row[0] for row in string_rows] == ["SRID=4326;POINT(1 2)", None]

binary_type = pa.struct(
[
pa.field("srid", pa.int32(), nullable=False),
pa.field("wkb", pa.binary(), nullable=False),
]
)
binary_schema = pa.schema([pa.field("g", binary_type, metadata=geo_metadata)])
binary_batch = pa.RecordBatch.from_arrays(
[pa.array([{"srid": 4326, "wkb": wkb}, None], type=binary_type)],
schema=binary_schema,
)
binary_rs = _make_rs(
_FakeKernelHandle(binary_schema, [binary_batch]), disable_pandas=True
)
assert binary_rs.description[0][1] == "geometry"
binary_rows = binary_rs.fetchall()
assert binary_rows[0][0] == {"srid": 4326, "wkb": wkb}
assert binary_rows[1][0] is None


def test_fetchmany_arrow_slices_within_batch(int_schema):
handle = _FakeKernelHandle(int_schema, [_batch(int_schema, [10, 20, 30, 40])])
rs = _make_rs(handle)
Expand Down
24 changes: 24 additions & 0 deletions tests/unit/test_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -476,6 +476,7 @@ def test_retry_and_socket_timeout_threaded_into_kernel_client(self):
_retry_stop_after_attempts_duration=600.0,
_socket_timeout=12.5,
_pool_maxsize=41,
enable_geospatial_support=False,
)
try:
_, kwargs = mock_kernel_client.call_args
Expand All @@ -486,6 +487,7 @@ def test_retry_and_socket_timeout_threaded_into_kernel_client(self):
assert opts["retry_stop_after_attempts_duration"] == 600.0
assert kwargs["request_timeout_secs"] == 12.5
assert kwargs["max_connections"] == 41
assert kwargs["enable_geospatial_support"] is False
finally:
conn.close()

Expand Down Expand Up @@ -791,6 +793,28 @@ def test_connect_use_kernel_instantiates_real_kernel_backend(self):
finally:
conn.close()

def test_enable_geospatial_support_matches_real_kernel_signature(self):
self._real_kernel_or_skip()

from databricks.sql.backend.kernel.client import (
_kernel_geospatial_kwargs,
_kernel_session_accepts_kwarg,
)
from databricks.sql.exc import NotSupportedError

# The ordinary kernel unit-test tier installs the latest published
# wheel, which may lag the KERNEL_REV source pin while the matching
# kernel change is still in flight. Validate the compatibility error
# in that case; once the wheel includes the typed option, validate
# both values against its real PyO3 signature.
if not _kernel_session_accepts_kwarg("enable_geospatial_support"):
with pytest.raises(NotSupportedError, match="newer databricks-sql-kernel"):
_kernel_geospatial_kwargs(True)
return

assert _kernel_geospatial_kwargs(True) == {"enable_geospatial_support": True}
assert _kernel_geospatial_kwargs(False) == {"enable_geospatial_support": False}


class TestReydenThriftFallback:
"""Transparent auto-recovery from a Reyden / Real-Time warehouse rejecting
Expand Down
Loading