Akash3121 commented on code in PR #10165:
URL: https://github.com/apache/paimon/pull/10165#discussion_r4097047425
##########
paimon-python/pypaimon/write/file_store_write.py:
##########
@@ -285,6 +280,11 @@ def _build_pk_merge_function(self):
# for the engines we know are out of scope today; any other
# NotImplementedError is a bug we want to surface, not swallow.
if engine == MergeEngine.AGGREGATE:
+ from pypaimon.read.merge_engine_support import check_supported
+
+ # Reject unsupported table options before any data is buffered.
+ # A read-side error cannot undo an incorrectly merged write.
+ check_supported(self.table)
Review Comment:
This check is not early enough to guarantee a side-effect-free rejection.
`TableWrite.write_arrow_batch` calls
`row_key_extractor.extract_partition_bucket_groups(data)` before it enters
`FileStoreWrite.write` and constructs the PK data writer. For `bucket=-1`,
`DynamicBucketRowKeyExtractor._extract_buckets_batch` records every new key
through `DynamicBucketIndexMaintainer.notify_new_record` before this line
raises. If the caller catches the write error and invokes `prepare_commit` - as
the new stream regression does - the index maintainer can write an index file
and return index-only commit messages even though no data row was accepted.
Please run aggregation validation at writer construction or before row-key
extraction, and add a dynamic-bucket regression asserting that a rejected write
leaves `prepare_commit` empty and creates no HASH index files.
Why it matters? the implementation fixes data buffering but can still leave
persistent hash-index state representing rows that were rejected.
##########
paimon-python/pypaimon/tests/test_aggregation_e2e.py:
##########
@@ -256,6 +264,34 @@ def test_field_ignore_retract_rejected(self):
'fields.total.ignore-retract',
)
+ def test_unsupported_stream_write_rejected_before_buffering(self):
+ table = self._create_pk_table(
+ 'agg_stream_reject', field_aggs={'total': 'sum'},
+ extra_options={'aggregation.remove-record-on-delete': 'true'})
+ writer = table.new_stream_write_builder().new_write()
+ try:
+ with self.assertRaisesRegex(
+ NotImplementedError,
'aggregation.remove-record-on-delete'):
+ writer.write_arrow(pa.Table.from_pylist([
+ {'id': 1, 'total': 10}, {'id': 1, 'total': 20},
+ ], schema=self.pa_schema))
+ self.assertEqual(writer.prepare_commit(1), [])
+ finally:
+ writer.close()
+ self.assertIsNone(table.snapshot_manager().get_latest_snapshot())
+ self.assertEqual(glob.glob(
+ os.path.join(table.table_path, '**', '*.parquet'),
recursive=True), [])
+
+ def test_false_retract_options_remain_writable(self):
+ table = self._create_pk_table(
+ 'agg_false_retract', field_aggs={'total': 'sum'}, extra_options={
+ 'aggregation.remove-record-on-delete': 'false',
+ 'fields.total.ignore-retract': 'false',
+ })
+ self._write(table, [{'id': 1, 'total': 10}])
Review Comment:
This new regression currently fails the required `Python` / `Native CI` job.
The writes complete, but the native reader rejects the table merely because the
false-valued keys are present: `NotImplementedError: ... options not supported
by this build: aggregation.remove-record-on-delete,
fields.total.ignore-retract`. That contradicts this PR’s stated guarantee that
explicitly false retract flags remain accepted and leaves CI red.
I think it would be better to either make the native validation use the same
truthiness semantics as `aggregation_unsupported_options`, or scope this test
to write acceptance only if native-read compatibility is intentionally outside
this change.
--
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]