This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 1fe6571bda [python] Add a table rollback CLI command (#10177)
1fe6571bda is described below
commit 1fe6571bda689db3634d19eb23c54f99458892ae
Author: jackylee <[email protected]>
AuthorDate: Fri Sep 25 21:09:11 2026 +0800
[python] Add a table rollback CLI command (#10177)
---
paimon-python/pypaimon/cli/cli_table.py | 78 +++++++++++++++++++++++++-
paimon-python/pypaimon/tests/cli_table_test.py | 51 +++++++++++++++++
2 files changed, 128 insertions(+), 1 deletion(-)
diff --git a/paimon-python/pypaimon/cli/cli_table.py
b/paimon-python/pypaimon/cli/cli_table.py
index 84b2b91d05..dc6f9d6481 100644
--- a/paimon-python/pypaimon/cli/cli_table.py
+++ b/paimon-python/pypaimon/cli/cli_table.py
@@ -413,6 +413,59 @@ def cmd_table_snapshot(args):
sys.exit(1)
+def cmd_table_rollback(args):
+ """
+ Execute the 'table rollback' command.
+
+ Rolls a table back to an earlier snapshot, tag, or timestamp using the
+ table's existing rollback API. Exactly one of --snapshot / --tag /
+ --timestamp must be provided.
+
+ Args:
+ args: Parsed command line arguments.
+ """
+ from pypaimon.cli.cli import load_catalog_config, create_catalog
+ from pypaimon.table.file_store_table import FileStoreTable
+
+ config = load_catalog_config(args.config)
+ catalog = create_catalog(config)
+
+ table_identifier = args.table
+ parts = table_identifier.split('.')
+ if len(parts) != 2:
+ print(f"Error: Invalid table identifier '{table_identifier}'. "
+ f"Expected format: 'database.table'", file=sys.stderr)
+ sys.exit(1)
+
+ database_name, table_name = parts
+ try:
+ table = catalog.get_table(f"{database_name}.{table_name}")
+ except Exception as e:
+ print(f"Error: Failed to get table '{table_identifier}': {e}",
file=sys.stderr)
+ sys.exit(1)
+
+ if not isinstance(table, FileStoreTable):
+ print(f"Error: Table '{table_identifier}' is not a FileStoreTable. "
+ f"Rollback operation is not supported for this table type.",
file=sys.stderr)
+ sys.exit(1)
+
+ try:
+ if args.snapshot is not None:
+ table.rollback_to(args.snapshot)
+ target_desc = f"snapshot {args.snapshot}"
+ elif args.tag is not None:
+ table.rollback_to(args.tag)
+ target_desc = f"tag '{args.tag}'"
+ else:
+ table.rollback_to_timestamp(args.timestamp)
+ target_desc = f"timestamp {args.timestamp}"
+ except Exception as e:
+ print(f"Error: Failed to roll back table '{table_identifier}': {e}",
file=sys.stderr)
+ sys.exit(1)
+
+ print(f"Rolled back table '{table_identifier}' to {target_desc}.")
+
+
def cmd_table_create(args):
"""
Execute the 'table create' command.
@@ -976,7 +1029,30 @@ def add_table_subcommands(table_parser):
help='Table identifier in format: database.table'
)
snapshot_parser.set_defaults(func=cmd_table_snapshot)
-
+
+ # table rollback command
+ rollback_parser = table_subparsers.add_parser(
+ 'rollback',
+ help='Roll back a table to an earlier snapshot, tag, or timestamp')
+ rollback_parser.add_argument(
+ 'table',
+ help='Table identifier in format: database.table'
+ )
+ rollback_target =
rollback_parser.add_mutually_exclusive_group(required=True)
+ rollback_target.add_argument(
+ '--snapshot', type=int, default=None,
+ help='Snapshot ID to roll back to'
+ )
+ rollback_target.add_argument(
+ '--tag', type=str, default=None,
+ help='Tag name to roll back to'
+ )
+ rollback_target.add_argument(
+ '--timestamp', type=int, default=None,
+ help='Epoch milliseconds; roll back to the latest snapshot at or
before it'
+ )
+ rollback_parser.set_defaults(func=cmd_table_rollback)
+
# table create command
create_parser = table_subparsers.add_parser('create', help='Create a new
table')
create_parser.add_argument(
diff --git a/paimon-python/pypaimon/tests/cli_table_test.py
b/paimon-python/pypaimon/tests/cli_table_test.py
index 95282ccf4b..1752e07397 100644
--- a/paimon-python/pypaimon/tests/cli_table_test.py
+++ b/paimon-python/pypaimon/tests/cli_table_test.py
@@ -1578,6 +1578,57 @@ class CliTableTest(unittest.TestCase):
self.assertEqual(ctx.exception.code, 1)
self.assertIn("Invalid WHERE clause", mock_stderr.getvalue())
+ def test_cli_table_rollback_to_snapshot(self):
+ """table rollback --snapshot restores an earlier snapshot."""
+ pa_schema = pa.schema([('id', pa.int32()), ('name', pa.string())])
+ schema = Schema.from_pyarrow_schema(pa_schema)
+ self.catalog.create_table('test_db.rollback_users', schema, True)
+
+ def _write(row_id, name):
+ table = self.catalog.get_table('test_db.rollback_users')
+ wb = table.new_batch_write_builder()
+ w = wb.new_write()
+ c = wb.new_commit()
+ w.write_arrow(pa.Table.from_pydict(
+ {'id': [row_id], 'name': [name]}, schema=pa_schema))
+ c.commit(w.prepare_commit())
+ w.close()
+ c.close()
+
+ _write(1, 'a') # snapshot 1
+ _write(2, 'b') # snapshot 2
+
+ def _latest_id():
+ fresh = CatalogFactory.create({'warehouse': self.warehouse})
+ return fresh.get_table(
+
'test_db.rollback_users').snapshot_manager().get_latest_snapshot().id
+
+ self.assertEqual(_latest_id(), 2)
+
+ with patch('sys.argv',
+ ['paimon', '-c', self.config_file, 'table', 'rollback',
+ 'test_db.rollback_users', '--snapshot', '1']):
+ with patch('sys.stdout', new_callable=StringIO) as out:
+ try:
+ main()
+ except SystemExit:
+ pass
+ self.assertIn('Rolled back', out.getvalue())
+
+ # The latest snapshot is back to 1 (snapshot 2 was rolled back).
+ self.assertEqual(_latest_id(), 1)
+
+ def test_cli_table_rollback_requires_a_target(self):
+ """table rollback without any target selector exits with an error."""
+ with patch('sys.argv',
+ ['paimon', '-c', self.config_file, 'table', 'rollback',
+ 'test_db.users']):
+ with patch('sys.stderr', new_callable=StringIO):
+ with self.assertRaises(SystemExit) as ctx:
+ main()
+ # argparse rejects the missing mutually-exclusive group.
+ self.assertNotEqual(ctx.exception.code, 0)
+
if __name__ == '__main__':
unittest.main()