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]

Reply via email to