Akash3121 commented on code in PR #10128:
URL: https://github.com/apache/paimon/pull/10128#discussion_r4078044604
##########
paimon-python/pypaimon/tests/reader_append_only_test.py:
##########
@@ -104,6 +104,37 @@ def test_avro_ao_reader(self):
actual = self._read_test_table(read_builder).sort_by('user_id')
self.assertEqual(actual, self.expected)
+ def test_avro_ao_reader_filter_keeps_rows_past_first_batch(self):
+ # A filtered Avro read must not stop when a full batch_size (1024)
block
+ # matches nothing: the reader returned None there, which the caller
reads as
+ # end-of-input, silently dropping every later matching row.
+ schema = Schema.from_pyarrow_schema(
+ self.pa_schema, partition_keys=['dt'], options={'file.format':
'avro'})
+ self.catalog.create_table('default.test_avro_filter_batches', schema,
False)
+ table = self.catalog.get_table('default.test_avro_filter_batches')
+
+ n = 3000
+ data = pa.Table.from_pydict({
+ 'user_id': list(range(n)),
+ 'item_id': [1000 + i for i in range(n)],
+ 'behavior': ['a'] * n,
+ 'dt': ['p1'] * n,
+ }, schema=self.pa_schema)
+ wb = table.new_batch_write_builder()
+ w, c = wb.new_write(), wb.new_commit()
+ w.write_arrow(data)
+ c.commit(w.prepare_commit())
+ w.close()
+ c.close()
+
+ pb = table.new_read_builder().new_predicate_builder()
+ read_builder = table.new_read_builder().with_filter(
+ pb.greater_or_equal('user_id', 2000))
Review Comment:
Nit: With the default batch size of 1024, this threshold makes only the
first batch fully filtered; the second batch already contains matching rows.
Since the implementation intentionally uses a loop to support arbitrarily many
consecutive empty batches, could we use 2500 here instead? With 3000 input
rows, that would skip two complete batches before returning the final 500 rows
and would distinguish the loop from an implementation that retries only once.
I think, the assertion should become:
```python
read_builder = table.new_read_builder().with_filter(
pb.greater_or_equal('user_id', 2500))
result = self._read_test_table(read_builder)
self.assertEqual(
sorted(result.column('user_id').to_pylist()), list(range(2500, n)))
```
--
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]