FIX: Implementing RecordBatchReader.Close() functionality for Arrow - #644
Conversation
📊 Code Coverage Report
Diff CoverageDiff: main...HEAD, staged and unstaged changes
Summary
mssql_python/cursor.pyLines 272-280 272 cursor = self._cursor
273 if cursor is not None and not cursor.closed and cursor.hstmt is not None:
274 try:
275 cursor.hstmt._cancel() # pylint: disable=protected-access
! 276 except Exception as e: # pylint: disable=broad-exception-caught
277 logger.debug("arrow_reader.close: SQLCancel raised: %s", e)
278
279 # Close the generator — this raises GeneratorExit inside it, which
280 # runs the try/finally cleanup block (SQLFreeStmt + diag drain +mssql_python/pybind/ddbc_bindings.cppLines 1445-1453 1445 SQLEndTran_ptr = GetFunctionPointer<SQLEndTranFunc>(handle, "SQLEndTran");
1446 SQLDisconnect_ptr = GetFunctionPointer<SQLDisconnectFunc>(handle, "SQLDisconnect");
1447 SQLFreeHandle_ptr = GetFunctionPointer<SQLFreeHandleFunc>(handle, "SQLFreeHandle");
1448 SQLFreeStmt_ptr = GetFunctionPointer<SQLFreeStmtFunc>(handle, "SQLFreeStmt");
! 1449 SQLCancel_ptr = GetFunctionPointer<SQLCancelFunc>(handle, "SQLCancel");
1450
1451 SQLGetDiagRec_ptr = GetFunctionPointer<SQLGetDiagRecFunc>(handle, "SQLGetDiagRecW");
1452
1453 SQLParamData_ptr = GetFunctionPointer<SQLParamDataFunc>(handle, "SQLParamData");Lines 1608-1617 1608 }
1609 }
1610
1611 void SqlHandle::cancel() {
! 1612 // SQLCancel is intentionally lenient: it is a no-op on non-STMT handles,
! 1613 // already-freed handles, or if the driver does not expose it. This lets
1614 // _ArrowReader.close() call it unconditionally without coordinating with
1615 // the fetch thread. The GIL is released so a blocked fetch thread can
1616 // observe the cancel and return.
1617 //Lines 1619-1628 1619 // The only cross-thread pattern this driver blesses is exactly the one
1620 // ODBC blesses: cancel() may be called from a thread *other than* the
1621 // fetch thread to unblock an in-flight SQLFetch/SQLExecute on the same
1622 // HSTMT. Per the ODBC spec, SQLCancel (with the SQLGetDiagRec/Field
! 1623 // family) is the only entry point safe to call across threads on the
! 1624 // same statement handle. All other operations on a Cursor/SqlHandle
1625 // are single-owner: per DB API 2.0 and the Cursor thread-safety note
1626 // in cursor.py, callers must not share a Cursor for its lifecycle
1627 // operations (execute/fetch/close/free) across threads. Under that
1628 // contract, free() / close_cursor() / SQLFreeHandle can never be inLines 1642-1651 1642 if (!SQLCancel_ptr) {
1643 return;
1644 }
1645 SQLHANDLE h = _handle;
! 1646 SQLRETURN ret;
! 1647 {
1648 py::gil_scoped_release release;
1649 ret = SQLCancel_ptr(h);
1650 }
1651 // SQLCancel may return SQL_SUCCESS_WITH_INFO when there was nothing to📋 Files Needing Attention📉 Files with overall lowest coverage (click to expand)mssql_python.pybind.logger_bridge.cpp: 59.2%
mssql_python.pybind.ddbc_bindings.h: 59.9%
mssql_python.pybind.logger_bridge.hpp: 70.8%
mssql_python.pybind.ddbc_bindings.cpp: 76.2%
mssql_python.__init__.py: 77.3%
mssql_python.row.py: 77.6%
mssql_python.ddbc_bindings.py: 79.6%
mssql_python.pybind.connection.connection_pool.cpp: 82.7%
mssql_python.pybind.connection.connection.cpp: 83.7%
mssql_python.connection.py: 84.7%🔗 Quick Links
|
There was a problem hiding this comment.
Pull request overview
This PR improves Arrow streaming ergonomics by making Cursor.arrow_reader() return a RecordBatchReader-compatible wrapper whose .close() actually cancels in-flight fetches and releases server-side ODBC cursor resources, keeping the parent Cursor reusable.
Changes:
- Introduces
_ArrowReaderincursor.pyand updatesCursor.arrow_reader()to return it instead of a rawpyarrow.RecordBatchReader. - Adds an ODBC
SQLCancelbinding and exposes it viaSqlHandle._cancel()to support cross-thread cancellation during.close(). - Extends Arrow reader tests to validate close semantics, context-manager behavior, GC cleanup, and cross-thread cancellation.
Reviewed changes
Copilot reviewed 4 out of 4 changed files in this pull request and generated 5 comments.
| File | Description |
|---|---|
mssql_python/cursor.py |
Adds _ArrowReader and reworks arrow_reader() to drive robust cleanup/cancellation and cursor state reset. |
mssql_python/pybind/ddbc_bindings.h |
Declares the SQLCancel function pointer type and adds SqlHandle::cancel() API surface. |
mssql_python/pybind/ddbc_bindings.cpp |
Loads SQLCancel, implements SqlHandle::cancel() (releasing the GIL), and exposes _cancel to Python. |
tests/test_004_cursor_arrow.py |
Updates/expands tests to cover new reader wrapper behavior and resource release semantics. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
|
Huh, I didn't know that arrow readers supported a close method! I think that |
Thanks for the suggestion — I looked into it and unfortunately subclassing doesn't work cleanly here because pyarrow.RecordBatchReader is a Cython extension type, not a normal Python class. Specifically: RecordBatchReader.from_batches(...) is a Cython factory that hard-codes the return type to the base class; Sub.from_batches(...) returns a plain RecordBatchReader, so overridden close() / read_next_batch() are never called. The current composition/duck-typing shape is intentional for that reason. Happy to add arrow_c_stream to the wrapper so consumers that follow the Arrow PyCapsule Protocol (polars, duckdb, RecordBatchReader.from_stream, etc.) can consume it without any isinstance check — that gives us the "no duck-typing worries" property without the subclass constraints. |
Add __arrow_c_stream__ to the _ArrowReader wrapper so Arrow-aware
consumers (pyarrow.RecordBatchReader.from_stream, polars.from_arrow,
duckdb.from_arrow, ...) can accept it directly without falling back to
isinstance(x, pa.RecordBatchReader) duck-typing.
Subclassing pyarrow.RecordBatchReader (Cython extension type) is not a
viable alternative: from_batches ignores the subclass, __class__
reassignment is rejected on Cython types, instances cannot hold
arbitrary attributes, and a bare Sub.__new__(Sub) segfaults as soon as
any inherited method touches the unset internal C stream. The PyCapsule
Protocol is the modern, standards-based interop mechanism that gives us
the 'no duck-typing worries' property without those constraints.
- Delegates to self._inner.__arrow_c_stream__ (pyarrow >= 14).
- Raises ArrowInvalid('Reader is closed') post-close, matching the
read_next_batch / schema semantics.
- Fails explicitly with a clear message on pyarrow < 14 instead of
silently returning an invalid capsule.
- Extends the class docstring to advertise PyCapsule support and to
explain why subclassing was ruled out (for future reviewers).
- Adds two tests: PyCapsule Protocol round-trip via
pa.RecordBatchReader.from_stream (skipped on pyarrow < 14) and a
post-close guard.
Gaurav Sharma (bewithgaurav)
left a comment
There was a problem hiding this comment.
suggesting a better impl of the wrapper
…O from thread-safety comment - cursor.py: change Cursor.arrow_reader return annotation from "pyarrow.RecordBatchReader" to "_ArrowReader" (the actual returned type). Update docstring to say the returned object *is* the wrapper rather than "behaves like" one. - mssql_python.pyi: mirror the annotation change and add a minimal _ArrowReader class stub documenting the wrapper's explicitly-owned surface (closed, schema, read_next_batch, close, __arrow_c_stream__, iteration and context-manager dunders) plus __getattr__ for runtime delegation to the wrapped pyarrow.RecordBatchReader. - ddbc_bindings.cpp: reword the DevSkim-flagged "TODO" inside the SqlHandle::cancel cross-thread-invariant comment to "note". The reference still points at the same thread-safety comment block in cursor.py; no code change.
Add 9 targeted tests for the previously-uncovered defensive branches of the Arrow reader's cancel/close path so the PR's newly-added code is covered by pytest rather than only the happy paths. _ArrowReader coverage (lines 167, 196, 213, 231, 237-238): - __getattr__ refuses leading-underscore names to prevent __del__ recursion. - __enter__ on a closed reader raises ArrowInvalid. - __arrow_c_stream__ raises ArrowInvalid when the inner reader lacks the PyCapsule Protocol (pyarrow < 14 fallback path). - __del__ short-circuits when sys.is_finalizing() is True. - __del__ swallows any exception raised inside the finalizer body. arrow_reader() cleanup generator coverage (lines 3046, 3052, 3060, 3082, 3091): - Finally block early-returns when the parent Cursor is already closed. - Pre- and post-close DDBCSQLGetAllDiagRecords failures are swallowed. - hstmt._close_cursor() failure is swallowed and cleanup continues to the bookkeeping-reset step. - Cursor._clear_rownumber() failure is swallowed. The four cleanup tests use a dedicated fresh mssql_python.connect() so the shared module-scoped 'cursor' fixture cannot hit 'connection is busy with results for another command' from residual state left by an earlier test. No production code changes.
Work Item / Issue Reference
Summary
This pull request introduces a new
_ArrowReaderwrapper to enhance the behavior of thearrow_readermethod in theCursorclass. The new wrapper ensures that server-side resources are properly released when the reader is closed, supporting robust cleanup and allowing the parent cursor to remain usable. The tests are updated and extended to verify the improved semantics, including close behavior and context manager support.Enhancements to Arrow batch reading and resource management:
_ArrowReaderclass tocursor.py, which wraps apyarrow.RecordBatchReaderand implements an 8-step close sequence to properly release server-side resources, reset cursor state, and support idempotent and context-manager-based cleanup. The parentCursorremains usable after closing the reader.Cursor.arrow_readermethod to return an instance of_ArrowReaderinstead of a rawpyarrow.RecordBatchReader, ensuring that closing the reader stops fetching, releases the server-side cursor, and resets cursor state. [1] [2]Test improvements for Arrow reader behavior:
test_arrow_readertest to check for duck-typed compatibility withpyarrow.RecordBatchReader, reflecting the new wrapper class.test_arrow_reader_close_semanticsto verify that.close()stops fetching, marks the reader as closed, is idempotent, and leaves the parent cursor usable; andtest_arrow_reader_context_managerto verify that the reader is closed on context manager exit and the cursor remains usable.