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

discivigour pushed a commit to branch feat/blobReaderClose
in repository https://gitbox.apache.org/repos/asf/paimon.git

commit a4775651f9801fbb37d9f06560cb735b8ac73eb6
Author: umi <[email protected]>
AuthorDate: Mon Jul 13 20:51:30 2026 +0800

    fix(python): close exhausted blob fallback readers
---
 .../pypaimon/read/reader/concat_batch_reader.py    | 13 ++--
 paimon-python/pypaimon/tests/blob_test.py          | 86 ++++++++++++++++++++++
 2 files changed, 93 insertions(+), 6 deletions(-)

diff --git a/paimon-python/pypaimon/read/reader/concat_batch_reader.py 
b/paimon-python/pypaimon/read/reader/concat_batch_reader.py
index e45418a018..317d15170d 100644
--- a/paimon-python/pypaimon/read/reader/concat_batch_reader.py
+++ b/paimon-python/pypaimon/read/reader/concat_batch_reader.py
@@ -415,9 +415,8 @@ class BlobFallbackBatchReader(RecordBatchReader):
             row_id - self._deletion_vector_range.from_
         )
 
-    @staticmethod
     def _state_overlaps_batch(
-        state: _BlobFileState, batch_first: int, batch_last: int
+        self, state: _BlobFileState, batch_first: int, batch_last: int
     ) -> bool:
         selected_ranges = state.selected_ranges
         while (
@@ -428,10 +427,12 @@ class BlobFallbackBatchReader(RecordBatchReader):
                 selected_ranges[state.selected_range_index].count()
             )
             state.selected_range_index += 1
-        return (
-            state.selected_range_index < len(selected_ranges)
-            and selected_ranges[state.selected_range_index].from_ <= batch_last
-        )
+        if state.selected_range_index >= len(selected_ranges):
+            # Batch row ids only move forward. Once the last selected range is
+            # behind this batch, the reader can never be used again.
+            self._close_state_reader(state)
+            return False
+        return selected_ranges[state.selected_range_index].from_ <= batch_last
 
     def _read_blob_values(
         self, state: _BlobFileState, batch_row_ids: List[int]
diff --git a/paimon-python/pypaimon/tests/blob_test.py 
b/paimon-python/pypaimon/tests/blob_test.py
index 86c7b5daa0..9174b5bde8 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -437,6 +437,92 @@ class BlobTest(unittest.TestCase):
         self.assertTrue(created_by_file["old.blob"][0].closed)
         self.assertTrue(created_by_file["new.blob"][0].closed)
 
+    def test_blob_fallback_batch_reader_closes_exhausted_readers(self):
+        created_by_file = {}
+
+        class FakeBlobReader:
+            def __init__(self, file_path, offset):
+                self._file_io = None
+                self.file_path = file_path
+                self.blob_lengths = [20]
+                self.blob_offsets = [offset]
+                self._input_stream = None
+                self.closed = False
+
+            def close(self):
+                self.closed = True
+
+        def data_file(name, first_row_id, max_sequence_number):
+            return DataFileMeta(
+                file_name=name,
+                file_size=0,
+                row_count=1,
+                min_key=None,
+                max_key=None,
+                key_stats=None,
+                value_stats=None,
+                min_sequence_number=max_sequence_number,
+                max_sequence_number=max_sequence_number,
+                schema_id=0,
+                level=0,
+                extra_files=[],
+                first_row_id=first_row_id,
+                file_path=name,
+            )
+
+        def supplier(file_path, offset):
+            def create_reader():
+                reader = FakeBlobReader(file_path, offset)
+                created_by_file.setdefault(file_path, []).append(reader)
+                return reader
+
+            return create_reader
+
+        def descriptor_offset(batch):
+            value = batch.column(0)[0].as_py()
+            return BlobDescriptor.deserialize(value).offset
+
+        reader = BlobFallbackBatchReader(
+            [
+                (data_file("first.blob", 0, 1), supplier("first.blob", 0)),
+                (
+                    data_file("second.blob", 10, 2),
+                    supplier("second.blob", 100),
+                ),
+                (
+                    data_file("third.blob", 20, 3),
+                    supplier("third.blob", 200),
+                ),
+            ],
+            "picture",
+            pa.large_binary(),
+            blob_as_descriptor=True,
+            batch_size=1,
+        )
+
+        first = reader.read_arrow_batch()
+        self.assertEqual(4, descriptor_offset(first))
+        self.assertFalse(created_by_file["first.blob"][0].closed)
+        self.assertNotIn("second.blob", created_by_file)
+
+        second = reader.read_arrow_batch()
+        self.assertEqual(104, descriptor_offset(second))
+        self.assertTrue(created_by_file["first.blob"][0].closed)
+        self.assertFalse(created_by_file["second.blob"][0].closed)
+        self.assertNotIn("third.blob", created_by_file)
+
+        third = reader.read_arrow_batch()
+        self.assertEqual(204, descriptor_offset(third))
+        self.assertTrue(created_by_file["second.blob"][0].closed)
+        self.assertFalse(created_by_file["third.blob"][0].closed)
+        self.assertIsNone(reader.read_arrow_batch())
+
+        self.assertEqual(1, len(created_by_file["first.blob"]))
+        self.assertEqual(1, len(created_by_file["second.blob"]))
+        self.assertEqual(1, len(created_by_file["third.blob"]))
+        reader.close()
+        self.assertTrue(created_by_file["third.blob"][0].closed)
+
     def test_blob_data_interface_compliance(self):
         """Test that BlobData properly implements Blob interface."""
         test_data = b"interface test data"

Reply via email to