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()

Reply via email to