Akash3121 commented on code in PR #10162:
URL: https://github.com/apache/paimon/pull/10162#discussion_r4095946433
##########
paimon-python/pypaimon/write/file_store_write.py:
##########
@@ -285,21 +285,38 @@ 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:
- # Surface the silent semantic mismatch in logs: the file
- # will be PK-unique (better than the pre-PR multi-row
- # corruption), but any reader that honours the declared
- # engine will see wrong values. Users sharing tables
- # across writers especially need to see this.
- logger.warning(
- "merge-engine '%s' is not implemented on the pypaimon "
- "write path; falling back to deduplicate so the flushed "
- "file stays PK-unique. The file contents reflect "
- "deduplicate semantics (latest writer wins), not %s "
- "semantics. Any reader that interprets the file under "
- "the declared engine will return incorrect results. "
- "Avoid the pypaimon writer for tables on this engine.",
- engine.value, engine.value)
- return DeduplicateMergeFunction()
+ from pypaimon.read.merge_engine_support import \
+ aggregation_unsupported_options
+ unsupported = aggregation_unsupported_options(self.table)
+ if unsupported:
+ # Out-of-scope aggregation options (retract opt-ins,
+ # sequence groups, aggregators like rbm64): the read-side
+ # ``check_supported`` guard raises for these, so don't
+ # build an aggregator the write buffer can't honor. Fall
+ # back to dedupe (PK-unique) as before; the user still gets
+ # the explicit error when a reader is built.
+ logger.warning(
+ "merge-engine 'aggregation' is configured with options "
+ "pypaimon does not implement (%s); the write buffer "
+ "falls back to deduplicate. Reading this table raises "
+ "until the options are removed.",
+ ", ".join(sorted(unsupported)))
+ return DeduplicateMergeFunction()
+ # Supported aggregation: aggregate same-key rows in the write
+ # buffer with the same AggregateMergeFunction the read path
+ # uses, so a single write_arrow carrying duplicate keys is
+ # partially aggregated instead of silently deduped
+ # (latest-row-wins). Read and compaction re-aggregate across
+ # files, so this mirrors Java's per-buffer partial aggregation.
+ from pypaimon.read.reader.aggregation_merge_function import (
+ AggregateMergeFunction, build_field_aggregators)
+ agg_value_fields = self.table.table_schema.fields
+ return AggregateMergeFunction(
Review Comment:
This is not safe when the table configures `sequence.field` . The read path
feeds `AggregateMergeFunction` in user-sequence order through
`builtin_seq_comparator` , but the write buffer sorts same-key rows only by
generated `_SEQUENCE_NUMBER` , so this new partial aggregation permanently
folds them in arrival order. For example, with ` sequence.field=total` , two
rows in one `write_arrow - (total=100, label='hi')` followed by `(total=50,
label='lo')` - are written as the single `total=50, label='lo'` row,
although the existing sequence-field contract requires the `total=100` row to
win. Once collapsed, read-side sorting cannot recover the discarded row. Please
either order each same-key run with the configured sequence comparator before
calling this merge function, or reject aggregation writes with
`sequence.field` until that ordering is implemented. Add a one-batch regression
where the highest sequence value is written first.
I reproduced the failure directly through
`KeyValueDataWriter._merge_pending_by_pk` : the two rows above collapse to
`{'total': 50, 'label': 'lo'}` . The existing sequence-field E2E test uses
separate commits, so it does not exercise the new same-buffer folding path.
--
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]