FIX: Implementing RecordBatchReader.Close() functionality for Arrow - #644
FIX: Implementing RecordBatchReader.Close() functionality for Arrow#644subrata-ms wants to merge 8 commits into
Conversation
📊 Code Coverage Report
Diff CoverageDiff: main...HEAD, staged and unstaged changes
Summary
mssql_python/cursor.pyLines 136-144 136 @property
137 def schema(self):
138 """Schema of the record batches produced by this reader."""
139 if self._closed:
! 140 raise self._arrow_invalid("Reader is closed")
141 return self._inner.schema
142
143 def read_next_batch(self):
144 if self._closed:Lines 168-176 168 # versions do not implement the protocol. Fail explicitly rather
169 # than silently returning something invalid.
170 inner_export = getattr(self._inner, "__arrow_c_stream__", None)
171 if inner_export is None:
! 172 raise self._arrow_invalid(
173 "Arrow PyCapsule Protocol requires pyarrow>=14; "
174 "the installed pyarrow version does not expose "
175 "RecordBatchReader.__arrow_c_stream__."
176 )Lines 185-193 185 return self._inner.read_next_batch()
186
187 def __enter__(self):
188 if self._closed:
! 189 raise self._arrow_invalid("Reader is closed")
190 return self
191
192 def __exit__(self, exc_type, exc_val, exc_tb):
193 self.close()Lines 203-211 203 try:
204 import sys as _sys
205
206 if _sys.is_finalizing():
! 207 return
208 # Retry whenever the generator is still referenced — covers both
209 # "user never called close()" and "earlier close() raised before
210 # the generator was released".
211 if getattr(self, "_generator", None) is not None:Lines 209-218 209 # "user never called close()" and "earlier close() raised before
210 # the generator was released".
211 if getattr(self, "_generator", None) is not None:
212 self.close()
! 213 except Exception: # pylint: disable=broad-exception-caught
! 214 pass
215
216 # ── Close implementation ──────────────────────────────────────────────
217
218 def close(self) -> None:Lines 248-256 248 cursor = self._cursor
249 if cursor is not None and not cursor.closed and cursor.hstmt is not None:
250 try:
251 cursor.hstmt._cancel() # pylint: disable=protected-access
! 252 except Exception as e: # pylint: disable=broad-exception-caught
253 logger.debug("arrow_reader.close: SQLCancel raised: %s", e)
254
255 # Close the generator — this raises GeneratorExit inside it, which
256 # runs the try/finally cleanup block (SQLFreeStmt + diag drain +Lines 2975-2983 2975 # body. This is the single canonical cleanup site.
2976 cur = cursor_ref[0]
2977 cursor_ref[0] = None
2978 if cur is None or cur.closed or cur.hstmt is None:
! 2979 return
2980
2981 # 1) Drain diagnostics produced by the (possibly cancelled)
2982 # fetch *before* SQL_CLOSE so we don't lose them.
2983 try:Lines 2981-2989 2981 # 1) Drain diagnostics produced by the (possibly cancelled)
2982 # fetch *before* SQL_CLOSE so we don't lose them.
2983 try:
2984 cur.messages.extend(ddbc_bindings.DDBCSQLGetAllDiagRecords(cur.hstmt))
! 2985 except Exception as e: # pylint: disable=broad-exception-caught
2986 logger.debug("arrow_reader cleanup: pre-close diag drain failed: %s", e)
2987
2988 # 2) Release the server-side cursor & locks while keeping the
2989 # HSTMT and prepared plan intact, so the parent Cursor canLines 2989-2997 2989 # HSTMT and prepared plan intact, so the parent Cursor can
2990 # be re-executed.
2991 try:
2992 cur.hstmt._close_cursor() # pylint: disable=protected-access
! 2993 except Exception as e: # pylint: disable=broad-exception-caught
2994 logger.debug("arrow_reader cleanup: _close_cursor failed: %s", e)
2995
2996 # 3) Drain diagnostics produced by SQL_CLOSE itself. This
2997 # runs unconditionally because SQL_CLOSE can returnLines 2999-3007 2999 # warning records on the HSTMT diag stack; the previous
3000 # "only on failure" path would silently drop those.
3001 try:
3002 cur.messages.extend(ddbc_bindings.DDBCSQLGetAllDiagRecords(cur.hstmt))
! 3003 except Exception as e: # pylint: disable=broad-exception-caught
3004 logger.debug("arrow_reader cleanup: post-close diag drain failed: %s", e)
3005
3006 # 4) Reset cursor bookkeeping to a clean "no result set"
3007 # state. rowcount becomes -1 to signal that the priorLines 3008-3016 3008 # result is no longer meaningful.
3009 try:
3010 cur._clear_rownumber() # pylint: disable=protected-access
3011 cur.rowcount = -1
! 3012 except Exception as e: # pylint: disable=broad-exception-caught
3013 logger.debug("arrow_reader cleanup: bookkeeping reset failed: %s", e)
3014
3015 gen = batch_generator()
3016 inner = pyarrow.RecordBatchReader.from_batches(schema, gen)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 1620-1629 1620 if (!SQLCancel_ptr) {
1621 return;
1622 }
1623 SQLHANDLE h = _handle;
! 1624 SQLRETURN ret;
! 1625 {
1626 py::gil_scoped_release release;
1627 ret = SQLCancel_ptr(h);
1628 }
1629 // 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: 81.4%
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.
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.