leaves12138 commented on code in PR #10081:
URL: https://github.com/apache/paimon/pull/10081#discussion_r4070100740
##########
paimon-python/pypaimon/write/table_commit.py:
##########
@@ -95,13 +96,52 @@ def _commit(
"Committing table %s, %d non-empty messages",
self.table.identifier, len(non_empty_messages)
)
+ if snapshot_properties is None:
+ prepared = self._prepare_native_commit(non_empty_messages)
+ if prepared is not None:
+ native, messages = prepared
+ # Mutation is deliberately outside the fallback boundary:
+ # an exception can mean the snapshot was already published.
+ native.commit(commit_identifier, messages)
+ return
self.file_store_commit.commit(**commit_kwargs)
+ def _prepare_native_commit(self, messages):
+ if (not self.table.options.native_commit_enabled()
+ or self.overwrite_partition is not None
+ or self._commit_callbacks):
+ return None
+ try:
+ from pypaimon.write.native_commit import (
+ create_native_commit, native_messages_supported,
+ to_native_commit_messages)
+ if not native_messages_supported(self.table, messages):
+ return None
+ if self._native_commit is None:
+ self._native_commit = create_native_commit(self.table,
self.commit_user)
+ if self._native_commit is None:
+ return None
+ return self._native_commit, to_native_commit_messages(self.table,
messages)
+ except Exception as error:
+ # No native mutation has started. Preserve the normal Python path
+ # when the optional runtime, FileIO or wire bridge is unavailable.
+ logger.debug("Native commit preparation failed; using Python: %s",
error)
+ return None
+
def abort(self, commit_messages: List[CommitMessage]):
+ prepared = self._prepare_native_commit(commit_messages)
+ if prepared is not None:
+ native, messages = prepared
+ native.abort(messages)
+ return
Review Comment:
[P2] Preserve Python partition paths when aborting native messages
With `commit.native.enabled=true`, a Python writer on a table partitioned by
`pt` still writes values such as `a/b` and `a=b` to the unescaped Python
directories (`pt=a/b/bucket-0/` and `pt=a=b/bucket-0/`). The wire conversion
intentionally removes `DataFileMeta.file_path`, so native `abort()`
reconstructs the Java-escaped paths (`pt=a%2Fb/...` / `pt=a%3Db/...`) instead.
Missing-file deletion is ignored by native cleanup, and this early return then
skips the Python cleanup. As a result, `abort()` returns successfully but
leaves the uncommitted Parquet files behind.
Reproduced on this head with the #912 runtime by changing the partition in
`test_native_abort_removes_uncommitted_files` from `None` to either `a/b` or
`a=b`: the file remains with native commit enabled, while the Python abort
route deletes it with the option disabled. Please retain Python cleanup for
incompatible local paths before invoking native abort, or preserve their actual
locations in the native cleanup adapter, and add these partition values to the
regression test.
--
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]