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]

Reply via email to