This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 158d83cb51 [python] Reuse streams for coalesced blob reads (#9133)
158d83cb51 is described below

commit 158d83cb51d9cdfaac9fce56d67018e8f83eae73
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Aug 10 17:19:18 2026 +0800

    [python] Reuse streams for coalesced blob reads (#9133)
---
 paimon-python/pypaimon/common/file_io.py  | 183 +++++++++-
 paimon-python/pypaimon/tests/blob_test.py | 546 +++++++++++++++++++++++++++++-
 2 files changed, 703 insertions(+), 26 deletions(-)

diff --git a/paimon-python/pypaimon/common/file_io.py 
b/paimon-python/pypaimon/common/file_io.py
index bfde96dc90..10406849b3 100644
--- a/paimon-python/pypaimon/common/file_io.py
+++ b/paimon-python/pypaimon/common/file_io.py
@@ -27,6 +27,8 @@ import pyarrow.fs as pafs
 
 from pypaimon.common.options import Options
 
+_LOG = logging.getLogger(__name__)
+
 
 def supports_pread(stream) -> bool:
     """Check if the stream supports position-based reads (thread-safe I/O)."""
@@ -53,6 +55,9 @@ def pread(stream, length: int, offset: int) -> bytes:
 _COALESCE_GAP = 1 << 20
 _COALESCE_SPAN = 8 << 20
 _COALESCE_VIEW_MAX_RETAINED_AMPLIFICATION = 2.0
+# Bound per-object opens; 16 cuts them by 75% for default 64-range batches.
+_MAX_RANGE_LANES_PER_PATH = 16
+_RANGE_REQUEST_WEIGHT = 1 << 20
 
 
 def create_temp_path(path: str) -> str:
@@ -186,8 +191,9 @@ class FileIO(ABC):
                               max_gap=_COALESCE_GAP, max_span=_COALESCE_SPAN):
         """Read ``ranges`` (each ``None`` or ``(path, offset, length)``), 
returning
         bytes in the same order. Same-file nearby ranges are merged into one 
read
-        to cut round trips, then sliced; reads run on a thread pool. Negative
-        length (read to EOF) is read on its own, never merged.
+        to cut round trips, then sliced. Each worker lane reuses one exclusive
+        stream for consecutive spans of the same path. Negative length (read to
+        EOF) is read on its own, never merged.
 
         A failed read propagates and aborts the whole batch (unlike a per-row
         ``file.open()`` loop that fails one row at a time).
@@ -232,12 +238,131 @@ class FileIO(ABC):
                 coalescible.append((index, path, offset, length))
 
         spans = _coalesce_ranges(coalescible, max_gap, max_span)
+        tasks_by_path = {}
+        for span in spans:
+            tasks_by_path.setdefault(span[0], []).append(("span", span))
+        for singleton in singletons:
+            tasks_by_path.setdefault(singleton[1], []).append(
+                ("one", singleton))
+        task_count = sum(len(path_tasks)
+                         for path_tasks in tasks_by_path.values())
+        if task_count == 0:
+            return results
+
+        workers = max(1, min(parallelism, task_count))
+
+        def _task_weight(task):
+            kind, payload = task
+            length = payload[2] if kind == "span" else payload[3]
+            return _RANGE_REQUEST_WEIGHT + max(0, length)
+
+        lanes = [[] for _ in range(workers)]
+        lane_loads = [0] * workers
+        path_task_groups = list(tasks_by_path.values())
+        path_loads = [
+            sum(_task_weight(task) for task in path_tasks)
+            for path_tasks in path_task_groups
+        ]
+        total_load = sum(path_loads)
+        path_capacities = [
+            min(len(path_tasks), _MAX_RANGE_LANES_PER_PATH)
+            for path_tasks in path_task_groups
+        ]
+        path_lane_counts = [
+            min(
+                capacity,
+                max(
+                    1,
+                    (workers * path_load + total_load - 1) // total_load,
+                ),
+            )
+            for path_load, capacity in zip(path_loads, path_capacities)
+        ]
+        remaining_lanes = max(
+            0,
+            min(workers, sum(path_capacities)) - sum(path_lane_counts),
+        )
+        for _ in range(remaining_lanes):
+            candidates = [
+                index for index in range(len(path_task_groups))
+                if path_lane_counts[index] < path_capacities[index]
+            ]
+            if not candidates:
+                break
+            index = max(
+                candidates,
+                key=lambda value: (
+                    path_loads[value] / path_lane_counts[value]
+                ),
+            )
+            path_lane_counts[index] += 1
+
+        for path_tasks, path_lanes in zip(
+                path_task_groups, path_lane_counts):
+            selected = sorted(
+                range(workers), key=lane_loads.__getitem__)[:path_lanes]
+            for task in sorted(path_tasks, key=_task_weight, reverse=True):
+                lane = min(selected, key=lane_loads.__getitem__)
+                lanes[lane].append(task)
+                lane_loads[lane] += _task_weight(task)
+        lanes = [lane for lane in lanes if lane]
+
+        class _RangeLane:
+            def __init__(self, file_io):
+                self._file_io = file_io
+                self._path = None
+                self._stream = None
+                self._close_error = None
+
+            def _close_current(self):
+                stream = self._stream
+                self._stream = None
+                self._path = None
+                if stream is None:
+                    return None
+                try:
+                    stream.close()
+                except BaseException as error:
+                    if self._close_error is None:
+                        self._close_error = error
+                    return error
+                return None
 
-        def _run(task):
+            def _stream_for(self, path):
+                if self._stream is not None and self._path == path:
+                    return self._stream
+                close_error = self._close_current()
+                if close_error is not None:
+                    raise close_error
+                self._stream = self._file_io.new_input_stream(path)
+                self._path = path
+                return self._stream
+
+            def read(self, path, offset, length):
+                try:
+                    stream = self._stream_for(path)
+                    if length >= 0 and supports_pread(stream):
+                        return pread(stream, length, offset)
+                    stream.seek(offset)
+                    return (stream.read() if length < 0
+                            else stream.read(length))
+                except Exception as read_error:
+                    self._close_current()
+                    if self._close_error is not None:
+                        raise read_error
+                    return self._file_io.read_file_range(
+                        path, offset, length)
+
+            def close(self):
+                self._close_current()
+                if self._close_error is not None:
+                    raise self._close_error
+
+        def _run_task(reader, task):
             kind, payload = task
             if kind == "span":
                 path, span_off, span_len, members = payload
-                buf = self.read_file_range(path, span_off, span_len)
+                buf = reader.read(path, span_off, span_len)
                 if return_views:
                     buf = memoryview(buf)
                     useful = sum(length for _, _, length in members)
@@ -246,22 +371,48 @@ class FileIO(ABC):
                         or span_len <= useful * max_retained_amplification
                     )
                 for idx, off, length in members:
-                    s = off - span_off
-                    value = buf[s:s + length]
+                    start = off - span_off
+                    value = buf[start:start + length]
                     if return_views and not share_buffer:
                         value = memoryview(bytes(value))
                     results[idx] = value
             else:
-                idx, path, off, length = payload
-                result = self.read_file_range(path, off, length)
-                results[idx] = memoryview(result) if return_views else result
-
-        tasks = [("span", s) for s in spans] + [("one", g) for g in singletons]
-        if not tasks:
-            return results
-        workers = max(1, min(parallelism, len(tasks)))
-        with ThreadPoolExecutor(workers) as pool:
-            list(pool.map(_run, tasks))
+                idx, path, offset, length = payload
+                result = reader.read(path, offset, length)
+                results[idx] = (
+                    memoryview(result) if return_views else result)
+
+        def _run_lane(lane):
+            reader = _RangeLane(self)
+            read_error = None
+            try:
+                for task in lane:
+                    _run_task(reader, task)
+            except BaseException as error:
+                read_error = error
+            close_error = None
+            try:
+                reader.close()
+            except BaseException as error:
+                close_error = error
+            return read_error, close_error
+
+        with ThreadPoolExecutor(len(lanes)) as pool:
+            outcomes = list(pool.map(_run_lane, lanes))
+        read_error = next(
+            (error for error, _ in outcomes if error is not None), None)
+        close_error = next(
+            (error for _, error in outcomes if error is not None), None)
+        if read_error is not None:
+            if close_error is not None:
+                _LOG.warning(
+                    "Failed to close a range input stream",
+                    exc_info=(type(close_error), close_error,
+                              close_error.__traceback__),
+                )
+            raise read_error
+        if close_error is not None:
+            raise close_error
         return results
 
     def read_blobs_concurrent(self, blobs, parallelism):
diff --git a/paimon-python/pypaimon/tests/blob_test.py 
b/paimon-python/pypaimon/tests/blob_test.py
index e5a14677cc..a20f24dcfe 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -22,6 +22,8 @@ import os
 import shutil
 import struct
 import tempfile
+import threading
+import time
 import unittest
 import zlib
 from decimal import Decimal
@@ -2412,18 +2414,36 @@ class BlobEndToEndTest(unittest.TestCase):
         calls = []
         range_reads = []
         original_read = file_io.read_blobs_concurrent
-        original_range_read = file_io.read_file_range
+        original_open = file_io.new_input_stream
 
         def read_blobs_concurrent(blobs, parallelism):
             calls.append((list(blobs), parallelism))
             return original_read(blobs, parallelism)
 
-        def read_file_range(path, offset, length):
-            range_reads.append((path, offset, length))
-            return original_range_read(path, offset, length)
+        def new_input_stream(path):
+            stream = original_open(path)
+
+            class TrackingStream:
+                def read(self, length=-1):
+                    return stream.read(length)
+
+                def seek(self, offset, whence=0):
+                    return stream.seek(offset, whence)
+
+                def tell(self):
+                    return stream.tell()
+
+                def read_at(self, length, offset):
+                    range_reads.append((path, offset, length))
+                    return os.pread(stream.fileno(), length, offset)
+
+                def close(self):
+                    stream.close()
+
+            return TrackingStream()
 
         file_io.read_blobs_concurrent = read_blobs_concurrent
-        file_io.read_file_range = read_file_range
+        file_io.new_input_stream = new_input_stream
         reader = FormatBlobReader(
             file_io=file_io,
             file_path=blob_file_path,
@@ -3714,13 +3734,22 @@ class CoalesceRangesTest(unittest.TestCase):
                 output.write(data)
             file_io = FileIO.get(f"file://{tmp_dir}", {})
             reads = []
-            original_read = file_io.read_file_range
+            original_open = file_io.new_input_stream
 
-            def read_file_range(file_path, offset, length):
-                reads.append((file_path, offset, length))
-                return original_read(file_path, offset, length)
+            def new_input_stream(file_path):
+                stream = original_open(file_path)
 
-            file_io.read_file_range = read_file_range
+                class TrackingStream:
+                    def read_at(self, length, offset):
+                        reads.append((file_path, offset, length))
+                        return os.pread(stream.fileno(), length, offset)
+
+                    def close(self):
+                        stream.close()
+
+                return TrackingStream()
+
+            file_io.new_input_stream = new_input_stream
             got = file_io.read_ranges_coalesced_views(
                 [(path, 0, 10), (path, 1000, 10)],
                 parallelism=4,
@@ -3744,6 +3773,503 @@ class CoalesceRangesTest(unittest.TestCase):
             self.assertEqual(reads, [(path, 0, 1010)])
             self.assertIs(shared[0].obj, shared[1].obj)
 
+    def test_lane_stream_failure_reopens_range(self):
+        from pypaimon.common.file_io import FileIO
+
+        data = bytes(range(64))
+        with tempfile.TemporaryDirectory() as tmp_dir:
+            path = os.path.join(tmp_dir, "blob.bin")
+            with open(path, "wb") as output:
+                output.write(data)
+            file_io = FileIO.get(f"file://{tmp_dir}", {})
+            fallbacks = []
+
+            class FailingStream:
+                def read_at(self, length, offset):
+                    raise IOError("shared stream failed")
+
+                def close(self):
+                    pass
+
+            def read_file_range(file_path, offset, length):
+                fallbacks.append((file_path, offset, length))
+                return data[offset:offset + length]
+
+            file_io.new_input_stream = lambda _: FailingStream()
+            file_io.read_file_range = read_file_range
+            ranges = [(path, 0, 4), (path, 16, 4)]
+
+            self.assertEqual(
+                [data[0:4], data[16:20]],
+                file_io.read_ranges_coalesced(
+                    ranges, parallelism=2, max_gap=0),
+            )
+            self.assertEqual(2, len(fallbacks))
+
+    def test_fallback_reads_do_not_exceed_parallelism(self):
+        from pypaimon.common.file_io import FileIO
+
+        parallelism = 8
+        file_io = FileIO.get("file:///tmp", {})
+        barrier = threading.Barrier(parallelism)
+        lock = threading.Lock()
+        open_streams = 0
+        max_open_streams = 0
+
+        def opened():
+            nonlocal open_streams, max_open_streams
+            with lock:
+                open_streams += 1
+                max_open_streams = max(max_open_streams, open_streams)
+
+        def closed():
+            nonlocal open_streams
+            with lock:
+                open_streams -= 1
+
+        class FailingStream:
+            def __init__(self):
+                self.closed = False
+                opened()
+
+            def read_at(self, length, offset):
+                barrier.wait(timeout=5)
+                raise IOError("pooled read failed")
+
+            def close(self):
+                if not self.closed:
+                    self.closed = True
+                    closed()
+
+        def read_file_range(path, offset, length):
+            opened()
+            try:
+                time.sleep(0.01)
+                return b"ok"
+            finally:
+                closed()
+
+        file_io.new_input_stream = lambda _: FailingStream()
+        file_io.read_file_range = read_file_range
+        ranges = [
+            ("blob-%d" % index, 0, 2)
+            for index in range(parallelism)
+        ]
+
+        self.assertEqual(
+            [b"ok"] * parallelism,
+            file_io.read_ranges_coalesced(
+                ranges, parallelism=parallelism, max_gap=0),
+        )
+        self.assertLessEqual(max_open_streams, parallelism)
+        self.assertEqual(0, open_streams)
+        self.assertFalse(barrier.broken)
+
+    def test_known_and_unknown_lengths_share_exclusive_lane(self):
+        from pypaimon.common.file_io import FileIO
+
+        file_io = FileIO.get("file:///tmp", {})
+        operations = []
+        streams = []
+
+        class LaneStream:
+            def __init__(self):
+                self.position = 0
+                self.closed = False
+
+            def read_at(self, length, offset):
+                operations.append(("read_at", offset, length))
+                return b"known"
+
+            def seek(self, offset):
+                operations.append(("seek", offset))
+                self.position = offset
+
+            def read(self):
+                operations.append(("read", self.position))
+                return b"tail"
+
+            def close(self):
+                self.closed = True
+
+        def new_input_stream(_):
+            stream = LaneStream()
+            streams.append(stream)
+            return stream
+
+        file_io.new_input_stream = new_input_stream
+        file_io.read_file_range = lambda *args: self.fail(
+            "exclusive lane unexpectedly used fallback")
+
+        self.assertEqual(
+            [b"known", b"tail"],
+            file_io.read_ranges_coalesced(
+                [("blob", 0, 5), ("blob", 5, -1)],
+                parallelism=1,
+                max_gap=0,
+            ),
+        )
+        self.assertEqual([
+            ("read_at", 0, 5),
+            ("seek", 5),
+            ("read", 5),
+        ], operations)
+        self.assertEqual(1, len(streams))
+        self.assertTrue(streams[0].closed)
+
+    def test_non_positional_streams_are_exclusive(self):
+        from pypaimon.common.file_io import FileIO
+
+        data = bytes(range(128))
+        file_io = FileIO.get("file:///tmp", {})
+
+        class SerialStream:
+            def __init__(self):
+                self.position = 0
+                self.reading = False
+
+            def seek(self, position):
+                self.position = position
+
+            def read(self, length):
+                if self.reading:
+                    raise AssertionError("non-positional reads overlapped")
+                self.reading = True
+                try:
+                    time.sleep(0.001)
+                    result = data[self.position:self.position + length]
+                    self.position += len(result)
+                    return result
+                finally:
+                    self.reading = False
+
+            def close(self):
+                pass
+
+        streams = []
+
+        def new_input_stream(_):
+            stream = SerialStream()
+            streams.append(stream)
+            return stream
+
+        file_io.new_input_stream = new_input_stream
+        ranges = [("blob", i * 4, 2) for i in range(16)]
+
+        self.assertEqual(
+            [data[i * 4:i * 4 + 2] for i in range(16)],
+            file_io.read_ranges_coalesced(
+                ranges, parallelism=8, max_gap=0),
+        )
+        self.assertGreater(len(streams), 1)
+        self.assertLessEqual(len(streams), 8)
+
+    def test_same_path_reuses_bounded_exclusive_lanes(self):
+        from pypaimon.common.file_io import FileIO
+
+        data = bytes(range(256)) * 64
+        file_io = FileIO.get("file:///tmp", {})
+        streams = []
+
+        class PositionalStream:
+            def __init__(self):
+                self.reading = False
+                self.reads = 0
+                self.closed = False
+
+            def read_at(self, length, offset):
+                if self.reading:
+                    raise AssertionError("one stream was used concurrently")
+                self.reading = True
+                try:
+                    time.sleep(0.01)
+                    self.reads += 1
+                    return data[offset:offset + length]
+                finally:
+                    self.reading = False
+
+            def close(self):
+                self.closed = True
+
+        def new_input_stream(_):
+            stream = PositionalStream()
+            streams.append(stream)
+            return stream
+
+        file_io.new_input_stream = new_input_stream
+        ranges = [("blob", i * 128, 16) for i in range(64)]
+
+        self.assertEqual(
+            [data[offset:offset + length]
+             for _, offset, length in ranges],
+            file_io.read_ranges_coalesced(
+                ranges, parallelism=64, max_gap=0),
+        )
+        self.assertEqual(16, len(streams))
+        self.assertTrue(all(stream.reads == 4 for stream in streams))
+        self.assertEqual(64, sum(stream.reads for stream in streams))
+        self.assertTrue(all(stream.closed for stream in streams))
+
+    def test_same_path_lanes_balance_estimated_io(self):
+        from pypaimon.common.file_io import FileIO
+
+        large = 8 << 20
+        file_io = FileIO.get("file:///tmp", {})
+        streams = []
+
+        class PositionalStream:
+            def __init__(self):
+                self.lengths = []
+
+            def read_at(self, length, offset):
+                self.lengths.append(length)
+                return b"x"
+
+            def close(self):
+                pass
+
+        def new_input_stream(_):
+            stream = PositionalStream()
+            streams.append(stream)
+            return stream
+
+        file_io.new_input_stream = new_input_stream
+        ranges = []
+        offset = 0
+        for index in range(256):
+            length = large if index % 16 == 0 else 1
+            ranges.append(("blob", offset, length))
+            offset += length + 1
+
+        file_io.read_ranges_coalesced(
+            ranges, parallelism=64, max_gap=0)
+
+        self.assertEqual(16, len(streams))
+        self.assertEqual(
+            [1] * 16,
+            sorted(stream.lengths.count(large) for stream in streams),
+        )
+
+    def test_skewed_paths_redistribute_capped_lanes(self):
+        from collections import Counter
+
+        from pypaimon.common.file_io import FileIO
+
+        file_io = FileIO.get("file:///tmp", {})
+        streams = Counter()
+        lock = threading.Lock()
+
+        class PositionalStream:
+            def __init__(self, path):
+                self.path = path
+
+            def read_at(self, length, offset):
+                return self.path[0].encode() * length
+
+            def close(self):
+                pass
+
+        def new_input_stream(path):
+            with lock:
+                streams[path] += 1
+            return PositionalStream(path)
+
+        file_io.new_input_stream = new_input_stream
+        ranges = (
+            [("hot", index * 2, 1) for index in range(9900)]
+            + [("cold", index * 2, 1) for index in range(100)]
+        )
+
+        result = file_io.read_ranges_coalesced(
+            ranges, parallelism=64, max_gap=0)
+
+        self.assertEqual([b"h"] * 9900 + [b"c"] * 100, result)
+        self.assertEqual(Counter({"hot": 16, "cold": 16}), streams)
+
+    def test_path_memberships_can_exceed_worker_count(self):
+        from collections import Counter
+
+        from pypaimon.common.file_io import FileIO
+
+        file_io = FileIO.get("file:///tmp", {})
+        streams = Counter()
+        lock = threading.Lock()
+        active_streams = 0
+        max_active_streams = 0
+
+        class PositionalStream:
+            def __init__(self, path):
+                self.path = path
+                self.closed = False
+
+            def read_at(self, length, offset):
+                return self.path[0].encode() * length
+
+            def close(self):
+                nonlocal active_streams
+                if self.closed:
+                    return
+                self.closed = True
+                with lock:
+                    active_streams -= 1
+
+        def new_input_stream(path):
+            nonlocal active_streams, max_active_streams
+            with lock:
+                streams[path] += 1
+                active_streams += 1
+                max_active_streams = max(
+                    max_active_streams, active_streams)
+            return PositionalStream(path)
+
+        file_io.new_input_stream = new_input_stream
+        ranges = (
+            [("hot", index * 2, 1) for index in range(10000)]
+            + [("cold-%d" % index, 0, 1) for index in range(63)]
+        )
+
+        result = file_io.read_ranges_coalesced(
+            ranges, parallelism=64, max_gap=0)
+
+        self.assertEqual(10063, len(result))
+        self.assertEqual(16, streams["hot"])
+        self.assertLessEqual(max_active_streams, 64)
+        self.assertEqual(0, active_streams)
+
+    def test_stream_count_is_bounded_across_paths(self):
+        from pypaimon.common.file_io import FileIO
+
+        file_io = FileIO.get("file:///tmp", {})
+        lock = threading.Lock()
+        open_streams = 0
+        max_open_streams = 0
+        total_streams = 0
+
+        class PositionalStream:
+            def read_at(self, length, offset):
+                time.sleep(0.03 if offset == 0 else 0.001)
+                return bytes([offset]) * length
+
+            def close(self):
+                nonlocal open_streams
+                with lock:
+                    open_streams -= 1
+
+        def new_input_stream(_):
+            nonlocal open_streams, max_open_streams, total_streams
+            with lock:
+                open_streams += 1
+                total_streams += 1
+                max_open_streams = max(max_open_streams, open_streams)
+            return PositionalStream()
+
+        file_io.new_input_stream = new_input_stream
+        ranges = [
+            ("blob-%d" % path, offset, 4)
+            for path in range(4)
+            for offset in range(0, 32, 8)
+        ]
+
+        self.assertEqual(
+            [bytes([offset]) * length for _, offset, length in ranges],
+            file_io.read_ranges_coalesced(
+                ranges, parallelism=4, max_gap=0),
+        )
+        self.assertLessEqual(max_open_streams, 4)
+        self.assertEqual(4, total_streams)
+        self.assertEqual(0, open_streams)
+
+    def test_closes_all_streams_before_raising_close_error(self):
+        from pypaimon.common.file_io import FileIO
+
+        data = bytes(range(128))
+        file_io = FileIO.get("file:///tmp", {})
+        streams = []
+
+        class CloseStream:
+            def __init__(self, index):
+                self.index = index
+                self.closed = False
+
+            def read_at(self, length, offset):
+                time.sleep(0.01)
+                return data[offset:offset + length]
+
+            def close(self):
+                self.closed = True
+                if self.index == 0:
+                    raise IOError("first close failed")
+
+        def new_input_stream(_):
+            stream = CloseStream(len(streams))
+            streams.append(stream)
+            return stream
+
+        file_io.new_input_stream = new_input_stream
+        ranges = [("blob", offset, 4) for offset in range(0, 64, 8)]
+
+        with self.assertRaisesRegex(IOError, "first close failed"):
+            file_io.read_ranges_coalesced(
+                ranges, parallelism=4, max_gap=0)
+
+        self.assertGreater(len(streams), 1)
+        self.assertTrue(all(stream.closed for stream in streams))
+
+    def test_close_error_does_not_mask_read_error(self):
+        from pypaimon.common.file_io import FileIO
+
+        file_io = FileIO.get("file:///tmp", {})
+
+        class FailingStream:
+            def __init__(self, path):
+                self.path = path
+
+            def read_at(self, length, offset):
+                if self.path == "read-error":
+                    raise IOError("shared read failed")
+                return b"ok"
+
+            def close(self):
+                raise IOError("close failed")
+
+        def fail_fallback(path, offset, length):
+            raise IOError("fallback read failed")
+
+        file_io.new_input_stream = FailingStream
+        file_io.read_file_range = fail_fallback
+
+        with self.assertRaisesRegex(IOError, "shared read failed"):
+            file_io.read_ranges_coalesced(
+                [("close-error", 0, 2), ("read-error", 0, 2)],
+                parallelism=2,
+                max_gap=0,
+            )
+
+    def test_failed_stream_close_stops_before_fallback(self):
+        from pypaimon.common.file_io import FileIO
+
+        file_io = FileIO.get("file:///tmp", {})
+        fallbacks = []
+
+        class FailingStream:
+            def read_at(self, length, offset):
+                raise IOError("pooled read failed")
+
+            def close(self):
+                raise IOError("discard close failed")
+
+        def read_file_range(path, offset, length):
+            fallbacks.append((path, offset, length))
+            return b"ok"
+
+        file_io.new_input_stream = lambda _: FailingStream()
+        file_io.read_file_range = read_file_range
+
+        with self.assertRaisesRegex(IOError, "pooled read failed"):
+            file_io.read_ranges_coalesced(
+                [("blob", 0, 2)], parallelism=1, max_gap=0)
+        self.assertEqual([], fallbacks)
+
 
 class ReadFileRangeTest(unittest.TestCase):
     """read_file_range must accept length == -1 (read to EOF) -- the valid

Reply via email to