This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 5ea33f1d feat(python): add abort method to TableCommit for write
rollback (#612)
5ea33f1d is described below
commit 5ea33f1d20cda4e97991b3f954ace777567ccccb
Author: hechao <[email protected]>
AuthorDate: Sun Jul 26 19:40:32 2026 +0800
feat(python): add abort method to TableCommit for write rollback (#612)
---
bindings/python/src/write.rs | 90 ++++++++++++++++++++++++++-----------
bindings/python/tests/test_write.py | 71 +++++++++++++++++++++++++++++
2 files changed, 135 insertions(+), 26 deletions(-)
diff --git a/bindings/python/src/write.rs b/bindings/python/src/write.rs
index 2b604296..07f0d7d3 100644
--- a/bindings/python/src/write.rs
+++ b/bindings/python/src/write.rs
@@ -175,40 +175,78 @@ pub struct PyTableCommit {
commit_user: String,
}
+/// Collect and validate commit messages from a Python iterable, returning the
+/// inner Rust `CommitMessage` values. Shared by `commit` and `abort`.
+fn collect_and_validate_messages<'py>(
+ messages: &Bound<'py, PyAny>,
+ table_location: &str,
+ commit_user: &str,
+ method: &str,
+) -> PyResult<Vec<CommitMessage>> {
+ let mut inner_messages = Vec::new();
+ let iter = messages.try_iter().map_err(|_| {
+ PyTypeError::new_err(format!(
+ "{method}() expects a sequence of CommitMessage objects"
+ ))
+ })?;
+ for item in iter {
+ let item = item?;
+ let msg: PyRef<'py, PyCommitMessage> = item.extract().map_err(|_| {
+ PyTypeError::new_err(format!(
+ "{method}() expects a sequence of CommitMessage objects"
+ ))
+ })?;
+ if msg.table_location != table_location {
+ return Err(PyValueError::new_err(format!(
+ "commit message was prepared for a different table \
+ (message table '{}', committer table '{}')",
+ msg.table_location, table_location
+ )));
+ }
+ if msg.commit_user != commit_user {
+ return Err(PyValueError::new_err(
+ "commit message was prepared by a different WriteBuilder \
+ (writer and committer must come from the same \
+ table.new_write_builder() so they share one commit_user)"
+ .to_string(),
+ ));
+ }
+ inner_messages.push(msg.inner.clone());
+ }
+ Ok(inner_messages)
+}
+
#[pymethods]
impl PyTableCommit {
/// Commit the given commit messages. Empty input is a no-op success.
fn commit(&self, py: Python<'_>, messages: &Bound<'_, PyAny>) ->
PyResult<()> {
- let mut inner_messages = Vec::new();
- let iter = messages.try_iter().map_err(|_| {
- PyTypeError::new_err("commit() expects a sequence of CommitMessage
objects")
- })?;
- for item in iter {
- let item = item?;
- let msg: PyRef<PyCommitMessage> = item.extract().map_err(|_| {
- PyTypeError::new_err("commit() expects a sequence of
CommitMessage objects")
- })?;
- if msg.table_location != self.table_location {
- return Err(PyValueError::new_err(format!(
- "commit message was prepared for a different table \
- (message table '{}', committer table '{}')",
- msg.table_location, self.table_location
- )));
- }
- if msg.commit_user != self.commit_user {
- return Err(PyValueError::new_err(
- "commit message was prepared by a different WriteBuilder \
- (writer and committer must come from the same \
- table.new_write_builder() so they share one commit_user)"
- .to_string(),
- ));
- }
- inner_messages.push(msg.inner.clone());
- }
+ let inner_messages = collect_and_validate_messages(
+ messages,
+ &self.table_location,
+ &self.commit_user,
+ "commit",
+ )?;
let rt = runtime();
py.detach(|| rt.block_on(async {
self.inner.commit(inner_messages).await }))
.map_err(to_py_err)
}
+
+ /// Abort a prepared commit by deleting newly written data, changelog and
+ /// index files. Deletion is best-effort: missing files or storage errors
+ /// are silently ignored so abort cleanup never masks the original write
+ /// failure. After abort the data must not be committed — the files no
+ /// longer exist.
+ fn abort(&self, py: Python<'_>, messages: &Bound<'_, PyAny>) ->
PyResult<()> {
+ let inner_messages = collect_and_validate_messages(
+ messages,
+ &self.table_location,
+ &self.commit_user,
+ "abort",
+ )?;
+ let rt = runtime();
+ py.detach(|| rt.block_on(async {
self.inner.abort(&inner_messages).await }))
+ .map_err(to_py_err)
+ }
}
/// An opaque commit message produced by `prepare_commit`, consumed by
`commit`.
diff --git a/bindings/python/tests/test_write.py
b/bindings/python/tests/test_write.py
index 8ef53e7a..aeac37d5 100644
--- a/bindings/python/tests/test_write.py
+++ b/bindings/python/tests/test_write.py
@@ -177,3 +177,74 @@ def test_commit_different_builder_same_table_raises():
messages = write.prepare_commit()
with pytest.raises(ValueError):
table.new_write_builder().new_commit().commit(messages)
+
+
+def test_abort_cleans_up_written_data():
+ # Write data, prepare commit, abort — the written files should be deleted
+ # and reading back should return zero rows.
+ with tempfile.TemporaryDirectory() as warehouse:
+ ctx = _make_empty_table(warehouse)
+ table = _get_table(warehouse)
+ wb = table.new_write_builder()
+ write = wb.new_write()
+ write.write_arrow(_batch([1, 2, 3], ["a", "b", "c"]))
+ messages = write.prepare_commit()
+ assert len(messages) >= 1
+ wb.new_commit().abort(messages)
+ # After abort, no snapshot was committed — the table should be empty.
+ batches = ctx.sql("SELECT COUNT(*) AS cnt FROM paimon.wdb.t")
+ assert batches[0].column(0).to_pylist() == [0]
+
+
+def test_abort_empty_messages_noop():
+ with tempfile.TemporaryDirectory() as warehouse:
+ ctx = _make_empty_table(warehouse)
+ table = _get_table(warehouse)
+ wb = table.new_write_builder()
+ messages = wb.new_write().prepare_commit() # no write
+ assert messages == []
+ wb.new_commit().abort(messages) # no-op success
+ batches = ctx.sql("SELECT COUNT(*) AS cnt FROM paimon.wdb.t")
+ assert batches[0].column(0).to_pylist() == [0]
+
+
+def test_abort_non_message_raises_typeerror():
+ with tempfile.TemporaryDirectory() as warehouse:
+ _make_empty_table(warehouse)
+ table = _get_table(warehouse)
+ with pytest.raises(TypeError):
+ table.new_write_builder().new_commit().abort([object()])
+ with pytest.raises(TypeError):
+ table.new_write_builder().new_commit().abort(42)
+
+
+def test_abort_cross_table_messages_raises():
+ with tempfile.TemporaryDirectory() as warehouse:
+ ctx = SQLContext()
+ ctx.register_catalog("paimon", {"warehouse": warehouse})
+ ctx.sql("CREATE SCHEMA paimon.wdb")
+ ctx.sql("CREATE TABLE paimon.wdb.t1 (id INT, name STRING)")
+ ctx.sql("CREATE TABLE paimon.wdb.t2 (id INT, name STRING)")
+ catalog = PaimonCatalog({"warehouse": warehouse})
+ t1 = catalog.get_table("wdb.t1")
+ t2 = catalog.get_table("wdb.t2")
+ batch = pa.record_batch(
+ [pa.array([1], pa.int32()), pa.array(["a"], pa.string())],
+ names=["id", "name"],
+ )
+ w1 = t1.new_write_builder().new_write()
+ w1.write_arrow(batch)
+ messages = w1.prepare_commit()
+ with pytest.raises(ValueError):
+ t2.new_write_builder().new_commit().abort(messages)
+
+
+def test_abort_different_builder_same_table_raises():
+ with tempfile.TemporaryDirectory() as warehouse:
+ _make_empty_table(warehouse)
+ table = _get_table(warehouse)
+ write = table.new_write_builder().new_write()
+ write.write_arrow(_batch([1], ["a"]))
+ messages = write.prepare_commit()
+ with pytest.raises(ValueError):
+ table.new_write_builder().new_commit().abort(messages)