leaves12138 commented on code in PR #9996:
URL: https://github.com/apache/paimon/pull/9996#discussion_r4056384499
##########
paimon-python/pypaimon/read/table_read.py:
##########
@@ -461,6 +523,69 @@ def _native_batches_parallel(
return [batch for group_batches in results for batch in group_batches]
+ def _native_batches_parallel_streaming(self, readers):
+ """Stream ready batches with one buffered batch per Rust reader."""
+ stop = threading.Event()
+ results = queue.Queue(maxsize=len(readers))
+ capacities = [threading.Semaphore(1) for _ in readers]
+ executor = ThreadPoolExecutor(
+ max_workers=len(readers),
+ thread_name_prefix="pypaimon-native-read",
+ )
+
+ def put(index, item):
+ while not stop.is_set():
+ try:
+ results.put((index, item), timeout=0.1)
+ return True
+ except queue.Full:
+ pass
+ return False
+
+ def read_group(index, batches):
+ try:
+ iterator = iter(batches)
+ while not stop.is_set():
+ if not capacities[index].acquire(timeout=0.1):
+ continue
+ if stop.is_set():
+ capacities[index].release()
+ return
+ try:
+ batch = next(iterator)
+ except StopIteration:
+ put(index, ('done', None))
+ return
+ if not put(index, ('batch', batch)):
+ return
+ except Exception as error:
+ put(index, ('error', error))
+ finally:
+ close = getattr(batches, 'close', None)
+ if close is not None:
+ close()
+
+ futures = [
+ executor.submit(read_group, index, reader)
+ for index, reader in enumerate(readers)
+ ]
+ try:
+ completed = 0
+ while completed < len(readers):
+ index, (kind, value) = results.get()
+ capacities[index].release()
+ if kind == 'batch':
+ yield value
+ elif kind == 'done':
+ completed += 1
+ else:
+ raise value
+ finally:
+ stop.set()
Review Comment:
[P1] Propagate close() from the public batch reader to the worker generator
The executor cleanup here only runs when this Python generator is closed or
exhausted. The public `to_arrow_batch_reader` returns
`RecordBatchReader.from_batches(...)` directly, and on PyArrow 19.0.1 calling
that reader's `close()` does not close the underlying generator.
I reproduced this through the public entry point with two native reader
groups, each with more than one batch:
```python
reader = read.to_arrow_batch_reader(splits, parallelism=2)
reader.read_next_batch()
reader.close()
```
After `close()`, both native readers remain unclosed and both
`pypaimon-native-read_*` threads remain alive, waiting for their capacity
permits. In an isolated subprocess, interpreter shutdown then hangs because
ThreadPoolExecutor workers are joined before the still-referenced reader is
destroyed; my 4-second timeout terminates it. This does not require a slow or
blocked storage request. The serial control exits normally, and
`_to_managed_arrow_batch_reader` also closes both workers and exits normally in
the same test.
Could the public streaming API expose the same close-aware ownership as the
managed/from_stream path, so closing the returned reader always closes
`_convert_native_batches` and this generator? Please add an early-close
regression through the public returned reader. The current prefetch test closes
the internal generator directly, which bypasses the missing lifecycle
propagation.
--
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]