gabotorresruiz commented on code in PR #44604:
URL: https://github.com/apache/superset/pull/44604#discussion_r4160493677
##########
superset/sql/execution/executor.py:
##########
@@ -186,6 +187,65 @@ def build_statement_blocks(
return parsed_script, blocks
+class _LimitedCursor:
+ """Bound cumulative cursor reads without bypassing engine fetch
processing."""
+
+ def __init__(self, cursor: Any, limit: int) -> None:
+ """Wrap a cursor with a shared budget for all row-reading methods."""
+ self._cursor = cursor
+ self._remaining = limit
+
+ def __getattr__(self, name: str) -> Any:
+ """Delegate metadata and driver-specific methods to the real cursor."""
+ return getattr(self._cursor, name)
+
+ @property
+ def arraysize(self) -> int:
+ """Expose the driver's default fetch batch size."""
+ return self._cursor.arraysize
+
+ @arraysize.setter
+ def arraysize(self, value: int) -> None:
+ """Preserve engine-specific cursor batch-size configuration."""
+ self._cursor.arraysize = value
+
+ def fetchmany(self, size: int | None = None) -> list[Any]:
+ """Read no more than the remaining budget, including across batches."""
+ size = self.arraysize if size is None else size
+ size = max(0, min(size, self._remaining))
+ # Some drivers interpret zero as unbounded, so do not call them at all.
+ if not size:
+ return []
+ rows = self._cursor.fetchmany(size)
+ self._remaining -= len(rows)
+ return rows
+
+ def check_truncated(self, db_engine_spec: type[BaseEngineSpec]) -> bool:
+ """Probe one extra row using the engine's fetch and error handling."""
+ return self._remaining == 0 and bool(
+ db_engine_spec.fetch_data(_LimitedCursor(self._cursor, 1))
+ )
Review Comment:
This block worries me a bit. The probe runs after `rows` has already been
fully fetched, so a driver error on that one extra round trip discards a
materialised result set and fails a query that had its data in hand.
I verified it on this branch with a cursor that serves the five requested
rows and then raises on the next fetch:
- probe active: `QueryStatus.FAILED`, `sqlite error: transient read
failure`, no statements returned, two driver fetches.
- same cursor with `check_truncated` stubbed to `False`:
`QueryStatus.SUCCESS`, 5 rows, one driver fetch.
I can see `test_execute_truncation_probe_engine_errors` asserts exactly this
for `OSError`, and the Drill `RuntimeError` case shows where you drew the line,
so I read it as deliberate. What I am unsure about is the shape of the trade.
The trigger is "the result reached the cap", which is the normal case for a
capped query on a large table, so every such query against a network engine now
carries a new failure window in exchange for a boolean. Would degrading be
better here?
```python
def check_truncated(self, db_engine_spec: type[BaseEngineSpec]) -> bool:
"""Probe one extra row using the engine's fetch and error handling."""
if self._remaining != 0:
return False
try:
return bool(db_engine_spec.fetch_data(_LimitedCursor(self._cursor,
1)))
except Exception: # pylint: disable=broad-except
logger.warning("Truncation probe failed; reporting the result as
partial")
return True
```
That keeps the honest answer (the caller is told the result may be partial)
without losing rows the engine already handed over. An extra case in
`test_execute_truncation_probe_engine_errors` asserting the rows survive an
`OSError` with `truncated is True` would pin it. Or am I misunderstanding
something, and a failed probe really should fail the query?
##########
superset/db_engine_specs/bigquery.py:
##########
@@ -488,16 +494,15 @@ def fetch_data(cls, cursor: Any, limit: int | None =
None) -> list[tuple[Any, ..
if has_request_context():
g.bq_memory_limited = memory_limited
g.bq_memory_limited_row_count = len(data)
- return data
-
- except Exception: # pylint: disable=broad-except
- # Broad catch on purpose: any failure in the size-estimation /
- # progressive-fetch path (BigQuery DB-API errors, network or
- # auth timeouts mid-fetch, ``sys.getsizeof`` raising on an
- # unexpected cell type, or a future ``Row`` subclass we don't
- # know how to unwrap) must degrade gracefully to the parent's
- # straight fetch so the user still gets data.
- # Fallback to parent implementation
+ return FetchedRows(data, truncated=memory_limited)
+
+ except Exception as ex: # pylint: disable=broad-except
+ # A forward-only cursor cannot replay the consumed sample. Falling
+ # back after an estimation, second-batch, or EOF-probe failure
would
+ # silently discard it and return only the remaining rows as
success.
+ if first_batch:
+ raise cls.get_dbapi_mapped_exception(ex) from ex
Review Comment:
Just a small NIT, and I think this is the right call: falling back after a
consumed sample would have returned rows 1001 onward as a success. It reaches
further than the MCP contract though. `BigQueryEngineSpec.fetch_data` is also
what legacy SQL Lab calls at `superset/sql_lab.py:351` and what dataset column
discovery calls at `superset/connectors/sqla/utils.py:207`, so a BigQuery mid
fetch error that used to degrade to a plain `fetchall()` now surfaces as a
query error on those paths too. The `UPDATING.md` entry covers the limit
contract and the SQL Lab limit change but not this one; a sentence there would
reach operators who never open the MCP page.
##########
superset/sql/execution/executor.py:
##########
@@ -186,6 +187,65 @@ def build_statement_blocks(
return parsed_script, blocks
+class _LimitedCursor:
+ """Bound cumulative cursor reads without bypassing engine fetch
processing."""
+
+ def __init__(self, cursor: Any, limit: int) -> None:
+ """Wrap a cursor with a shared budget for all row-reading methods."""
+ self._cursor = cursor
+ self._remaining = limit
+
+ def __getattr__(self, name: str) -> Any:
+ """Delegate metadata and driver-specific methods to the real cursor."""
+ return getattr(self._cursor, name)
+
+ @property
+ def arraysize(self) -> int:
+ """Expose the driver's default fetch batch size."""
+ return self._cursor.arraysize
+
+ @arraysize.setter
+ def arraysize(self, value: int) -> None:
+ """Preserve engine-specific cursor batch-size configuration."""
+ self._cursor.arraysize = value
+
+ def fetchmany(self, size: int | None = None) -> list[Any]:
+ """Read no more than the remaining budget, including across batches."""
+ size = self.arraysize if size is None else size
+ size = max(0, min(size, self._remaining))
+ # Some drivers interpret zero as unbounded, so do not call them at all.
+ if not size:
+ return []
+ rows = self._cursor.fetchmany(size)
+ self._remaining -= len(rows)
+ return rows
+
+ def check_truncated(self, db_engine_spec: type[BaseEngineSpec]) -> bool:
+ """Probe one extra row using the engine's fetch and error handling."""
+ return self._remaining == 0 and bool(
+ db_engine_spec.fetch_data(_LimitedCursor(self._cursor, 1))
+ )
+
+ def fetchall(self) -> list[Any]:
+ """Translate an unbounded read into a bounded driver fetch."""
+ return self.fetchmany(self._remaining)
Review Comment:
Not a blocker, and it is the note I dropped last time, but the new flag
raises the stakes enough that I want to put it on the record. `fetchall()`
issues a single `fetchmany(self._remaining)`, and PEP 249 lets a driver return
fewer rows than `size` without being exhausted. When that happens the budget is
not consumed, so `check_truncated` short circuits to `False` and the response
now affirmatively reports a short result as complete.
Measured on this branch with a cursor that caps every batch at two rows, ten
rows available and a budget of five: `BaseEngineSpec.fetch_data` returns 2
rows, `_remaining` stays at 3, and `check_truncated` reports `False`. Before
this PR such a result was merely short; now it is short and labelled complete.
I still cannot name a driver in the tree that short batches, so this is
cheap insurance rather than a known break. Looping in `fetchall()` until the
budget is spent or a batch comes back empty would close it, and a
`test_limited_cursor_shares_read_budget` case with a short batching mock would
pin it.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]