From bdb9f65e72d12fba183d2804136d34d0381fbf35 Mon Sep 17 00:00:00 2001 From: "wenchao.wu" Date: Sat, 10 Oct 2026 10:21:50 +0800 Subject: [PATCH] [python] Support tag rename and replace in the CLI Filesystem catalogs rename and replace tag files. REST catalogs reject both because the API has no endpoint, so the CLI does not rewrite tag files behind the server. --- docs/docs/pypaimon/cli-versions.md | 30 ++++++- paimon-python/pypaimon/catalog/catalog.py | 34 ++++++++ .../pypaimon/catalog/filesystem_catalog.py | 49 +++++++++++ .../pypaimon/catalog/rest/rest_catalog.py | 18 ++++ paimon-python/pypaimon/cli/cli_tag.py | 83 ++++++++++++++++++- paimon-python/pypaimon/tests/cli_tag_test.py | 80 ++++++++++++++++++ .../pypaimon/tests/rest/rest_tag_test.py | 19 +++++ 7 files changed, 308 insertions(+), 5 deletions(-) diff --git a/docs/docs/pypaimon/cli-versions.md b/docs/docs/pypaimon/cli-versions.md index c64cae2f92f0..09df89ed0f96 100644 --- a/docs/docs/pypaimon/cli-versions.md +++ b/docs/docs/pypaimon/cli-versions.md @@ -32,7 +32,7 @@ Retain a dataset snapshot with a tag or work on a separate table history with a Manage tags (named snapshots) on a table. Tags are useful for time travel and pinning a snapshot for later access. ```shell -paimon tag mydb.users ... +paimon tag mydb.users ... ``` ### Tag Create @@ -46,11 +46,15 @@ paimon tag create mydb.users v1 --snapshot-id 3 # Do not error if the tag already exists paimon tag create mydb.users v1 --ignore-if-exists + +# Retain the tag for 12 hours +paimon tag create mydb.users v1 --time-retained 12h ``` Options: - `--snapshot-id, -s`: Snapshot id to tag (default: the latest snapshot) - `--ignore-if-exists, -i`: Do not raise an error if the tag already exists +- `--time-retained, -r`: Retention for the new tag, for example `1d` or `12h`. Omit it to store a plain snapshot reference. ### Tag List @@ -87,6 +91,30 @@ Options: paimon tag delete mydb.users v1 ``` +### Tag Rename + +```shell +paimon tag rename mydb.users v1 v2 +``` + +Filesystem catalogs rename the tag file. A REST catalog has no rename endpoint and rejects the command. + +### Tag Replace + +Point an existing tag at another snapshot. Without `--snapshot-id`, the tag follows the latest snapshot. `--time-retained` stores a create time and TTL on the tag. Omit it and the tag is rewritten as a plain snapshot reference, which drops a retention that was already set. A REST catalog has no replace endpoint and rejects the command. + +```shell +# Point v1 at the latest snapshot +paimon tag replace mydb.users v1 + +# Point v1 at snapshot 3 and retain it for 12 hours +paimon tag replace mydb.users v1 --snapshot-id 3 --time-retained 12h +``` + +Options: +- `--snapshot-id, -s`: Snapshot id to point at (default: the latest snapshot) +- `--time-retained, -r`: Retention for the replaced tag, for example `1d` or `12h` + ## Branch Commands Manage separate table histories. Create an empty branch with the current schema, diff --git a/paimon-python/pypaimon/catalog/catalog.py b/paimon-python/pypaimon/catalog/catalog.py index 958e84b05418..300d9c722780 100644 --- a/paimon-python/pypaimon/catalog/catalog.py +++ b/paimon-python/pypaimon/catalog/catalog.py @@ -425,5 +425,39 @@ def list_tags_paged( "list_tags_paged is not supported by this catalog." ) + def rename_tag( + self, + identifier: Union[str, Identifier], + tag_name: str, + target_tag_name: str, + ) -> None: + """Rename a tag on a table. + + Raises: + NotImplementedError: If the catalog does not support tag management. + """ + raise NotImplementedError( + "rename_tag is not supported by this catalog." + ) + + def replace_tag( + self, + identifier: Union[str, Identifier], + tag_name: str, + snapshot_id: Optional[int] = None, + time_retained: Optional[str] = None, + ) -> None: + """Point an existing tag at a snapshot. + + ``snapshot_id`` defaults to the latest snapshot. ``time_retained`` + is written only when set. + + Raises: + NotImplementedError: If the catalog does not support tag management. + """ + raise NotImplementedError( + "replace_tag is not supported by this catalog." + ) + def auth_table_query(self, identifier: Identifier, select: Optional[List[str]]) -> TableQueryAuthResult: raise NotImplementedError("auth_table_query not supported by this catalog") diff --git a/paimon-python/pypaimon/catalog/filesystem_catalog.py b/paimon-python/pypaimon/catalog/filesystem_catalog.py index d1264e2f9b43..5e169021ca43 100644 --- a/paimon-python/pypaimon/catalog/filesystem_catalog.py +++ b/paimon-python/pypaimon/catalog/filesystem_catalog.py @@ -445,6 +445,40 @@ def list_tags_paged( next_token = None return PagedList(elements=page, next_page_token=next_token) + def rename_tag( + self, + identifier: Union[str, Identifier], + tag_name: str, + target_tag_name: str, + ) -> None: + if not isinstance(identifier, Identifier): + identifier = Identifier.from_string(identifier) + table = self.get_table(identifier) + try: + table.rename_tag(tag_name, target_tag_name) + except ValueError as e: + _reraise_tag_value_error( + e, missing_tag=tag_name, conflict_tag=target_tag_name) + + def replace_tag( + self, + identifier: Union[str, Identifier], + tag_name: str, + snapshot_id: Optional[int] = None, + time_retained: Optional[str] = None, + ) -> None: + if not isinstance(identifier, Identifier): + identifier = Identifier.from_string(identifier) + table = self.get_table(identifier) + try: + table.replace_tag( + tag_name, + snapshot_id=snapshot_id, + time_retained=time_retained, + ) + except ValueError as e: + _reraise_tag_value_error(e, missing_tag=tag_name) + # ===================== Branch CRUD ===================== # Thin wrappers that delegate to FileSystemBranchManager (returned by # FileStoreTable.branch_manager() in the local-catalog case). Mirrors @@ -536,3 +570,18 @@ def list_branches( identifier = Identifier.from_string(identifier) table = self.get_table(identifier) return table.branch_manager().branches() + + +def _reraise_tag_value_error(exc, missing_tag, conflict_tag=None): + """Map TagManager ValueError text onto catalog exceptions. + + Snapshot-missing and blank-name errors stay ValueError. Only a missing + or conflicting tag name is translated. The ``Tag `` prefix keeps a + snapshot error such as ``Snapshot id '999' doesn't exist.`` as ValueError. + """ + message = str(exc) + if message.startswith("Tag ") and "doesn't exist" in message: + raise TagNotExistException(missing_tag) from exc + if conflict_tag is not None and "already exists" in message: + raise TagAlreadyExistException(conflict_tag) from exc + raise diff --git a/paimon-python/pypaimon/catalog/rest/rest_catalog.py b/paimon-python/pypaimon/catalog/rest/rest_catalog.py index ac33e55c3003..3425de48b045 100644 --- a/paimon-python/pypaimon/catalog/rest/rest_catalog.py +++ b/paimon-python/pypaimon/catalog/rest/rest_catalog.py @@ -624,6 +624,24 @@ def delete_tag(self, identifier: Union[str, Identifier], tag_name: str) -> None: except ForbiddenException as e: raise TableNoPermissionException(identifier) from e + def rename_tag(self, identifier: Union[str, Identifier], tag_name: str, + target_tag_name: str) -> None: + # The REST API has create, get, list, and delete only. Renaming the + # tag file through FileStoreTable would leave the server registry + # pointing at a name whose file is gone. + raise NotImplementedError( + "REST catalog does not support rename_tag. " + "The REST API has no rename endpoint." + ) + + def replace_tag(self, identifier: Union[str, Identifier], tag_name: str, + snapshot_id: Optional[int] = None, + time_retained: Optional[str] = None) -> None: + raise NotImplementedError( + "REST catalog does not support replace_tag. " + "The REST API has no replace endpoint." + ) + # Branch CRUD: mirrors Java RESTCatalog branch handlers. def create_branch(self, identifier: Union[str, Identifier], branch_name: str, tag_name: Optional[str] = None) -> None: diff --git a/paimon-python/pypaimon/cli/cli_tag.py b/paimon-python/pypaimon/cli/cli_tag.py index eb513ea606fe..91ab57f41560 100644 --- a/paimon-python/pypaimon/cli/cli_tag.py +++ b/paimon-python/pypaimon/cli/cli_tag.py @@ -17,10 +17,10 @@ """Tag commands for the Paimon CLI. -Adds the top-level ``tag {create,list,delete,get}`` subcommands (alongside -``table`` / ``db`` / ``catalog``). All operations go through the Catalog -layer so they get typed exceptions and work for both filesystem and REST -catalogs. +Adds the top-level ``tag {create,list,get,delete,rename,replace}`` +subcommands (alongside ``table`` / ``db`` / ``catalog``). All operations go +through the Catalog layer so they get typed exceptions and work for both +filesystem and REST catalogs. """ import json @@ -61,6 +61,7 @@ def cmd_tag_create(args): identifier, args.tag_name, snapshot_id=args.snapshot_id, + time_retained=args.time_retained, ignore_if_exists=args.ignore_if_exists, ) except TableNotExistException: @@ -96,6 +97,54 @@ def cmd_tag_delete(args): print("Tag '{}' deleted from table '{}'.".format(args.tag_name, identifier)) +def cmd_tag_rename(args): + """Execute ``tag rename``.""" + catalog, identifier = _open_catalog(args) + try: + catalog.rename_tag(identifier, args.tag_name, args.target_tag_name) + except TableNotExistException: + print("Error: Table '{}' does not exist.".format(identifier), + file=sys.stderr) + sys.exit(1) + except TagNotExistException: + print("Error: Tag '{}' does not exist.".format(args.tag_name), + file=sys.stderr) + sys.exit(1) + except TagAlreadyExistException: + print("Error: Tag '{}' already exists.".format(args.target_tag_name), + file=sys.stderr) + sys.exit(1) + except Exception as e: + print("Error: Failed to rename tag: {}".format(e), file=sys.stderr) + sys.exit(1) + print("Tag '{}' renamed to '{}' on table '{}'.".format( + args.tag_name, args.target_tag_name, identifier)) + + +def cmd_tag_replace(args): + """Execute ``tag replace``.""" + catalog, identifier = _open_catalog(args) + try: + catalog.replace_tag( + identifier, + args.tag_name, + snapshot_id=args.snapshot_id, + time_retained=args.time_retained, + ) + except TableNotExistException: + print("Error: Table '{}' does not exist.".format(identifier), + file=sys.stderr) + sys.exit(1) + except TagNotExistException: + print("Error: Tag '{}' does not exist.".format(args.tag_name), + file=sys.stderr) + sys.exit(1) + except Exception as e: + print("Error: Failed to replace tag: {}".format(e), file=sys.stderr) + sys.exit(1) + print("Tag '{}' replaced on table '{}'.".format(args.tag_name, identifier)) + + def cmd_tag_list(args): """Execute ``tag list``.""" catalog, identifier = _open_catalog(args) @@ -179,6 +228,9 @@ def add_tag_subcommands(subparsers): create_parser.add_argument( '--ignore-if-exists', '-i', action='store_true', help='Do not error if the tag already exists') + create_parser.add_argument( + '--time-retained', '-r', default=None, + help='Retention for the new tag, for example 1d or 12h') create_parser.set_defaults(func=cmd_tag_create) # tag list @@ -212,3 +264,26 @@ def add_tag_subcommands(subparsers): 'table', help='Table identifier in format: database.table') delete_parser.add_argument('tag_name', help='Name of the tag to delete') delete_parser.set_defaults(func=cmd_tag_delete) + + # tag rename + rename_parser = tag_subparsers.add_parser( + 'rename', help='Rename a tag') + rename_parser.add_argument( + 'table', help='Table identifier in format: database.table') + rename_parser.add_argument('tag_name', help='Current tag name') + rename_parser.add_argument('target_tag_name', help='New tag name') + rename_parser.set_defaults(func=cmd_tag_rename) + + # tag replace + replace_parser = tag_subparsers.add_parser( + 'replace', help='Point an existing tag at a snapshot') + replace_parser.add_argument( + 'table', help='Table identifier in format: database.table') + replace_parser.add_argument('tag_name', help='Name of the tag to replace') + replace_parser.add_argument( + '--snapshot-id', '-s', type=int, default=None, + help='Snapshot id to point at (default: the latest snapshot)') + replace_parser.add_argument( + '--time-retained', '-r', default=None, + help='Retention for the replaced tag, for example 1d or 12h') + replace_parser.set_defaults(func=cmd_tag_replace) diff --git a/paimon-python/pypaimon/tests/cli_tag_test.py b/paimon-python/pypaimon/tests/cli_tag_test.py index 9e3b75b9d17b..7e75d791163c 100644 --- a/paimon-python/pypaimon/tests/cli_tag_test.py +++ b/paimon-python/pypaimon/tests/cli_tag_test.py @@ -133,6 +133,13 @@ def test_create_duplicate_ignore_if_exists(self): 'tag', 'create', 'db.t', 'v1', '--ignore-if-exists') self.assertEqual(0, code) + def test_create_time_retained(self): + _, _, code = self._run( + 'tag', 'create', 'db.t', 'v1', '--time-retained', '12h') + self.assertEqual(0, code) + got, _, _ = self._run('tag', 'get', 'db.t', 'v1') + self.assertIn("Time Retained: PT12H", got) + # -- get ----------------------------------------------------------------- def test_get_table_format(self): @@ -170,6 +177,79 @@ def test_delete_not_exists(self): self.assertEqual(1, code) self.assertIn("does not exist", err) + # -- rename -------------------------------------------------------------- + + def test_rename(self): + self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1') + out, _, code = self._run('tag', 'rename', 'db.t', 'v1', 'v2') + self.assertEqual(0, code) + self.assertIn("renamed", out) + listed, _, _ = self._run('tag', 'list', 'db.t', '--format', 'json') + self.assertEqual(["v2"], json.loads(listed)) + got, _, _ = self._run('tag', 'get', 'db.t', 'v2') + self.assertIn("Snapshot ID: 1", got) + + def test_rename_missing(self): + _, err, code = self._run('tag', 'rename', 'db.t', 'absent', 'v2') + self.assertEqual(1, code) + self.assertIn("does not exist", err) + + def test_rename_target_exists(self): + self._run('tag', 'create', 'db.t', 'v1') + self._run('tag', 'create', 'db.t', 'v2') + _, err, code = self._run('tag', 'rename', 'db.t', 'v1', 'v2') + self.assertEqual(1, code) + self.assertIn("already exists", err) + + # -- replace ------------------------------------------------------------- + + def test_replace_snapshot(self): + self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1') + out, _, code = self._run( + 'tag', 'replace', 'db.t', 'v1', '--snapshot-id', '2') + self.assertEqual(0, code) + self.assertIn("replaced", out) + got, _, _ = self._run('tag', 'get', 'db.t', 'v1') + self.assertIn("Snapshot ID: 2", got) + + def test_replace_latest(self): + self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1') + _, _, code = self._run('tag', 'replace', 'db.t', 'v1') + self.assertEqual(0, code) + got, _, _ = self._run('tag', 'get', 'db.t', 'v1') + self.assertIn("Snapshot ID: 2", got) + + def test_replace_time_retained(self): + self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1') + _, _, code = self._run( + 'tag', 'replace', 'db.t', 'v1', '--time-retained', '12h') + self.assertEqual(0, code) + got, _, _ = self._run('tag', 'get', 'db.t', 'v1') + self.assertIn("Time Retained: PT12H", got) + + def test_replace_missing_tag(self): + _, err, code = self._run( + 'tag', 'replace', 'db.t', 'absent', '--snapshot-id', '1') + self.assertEqual(1, code) + self.assertIn("does not exist", err) + + def test_replace_missing_snapshot(self): + self._run('tag', 'create', 'db.t', 'v1') + _, err, code = self._run( + 'tag', 'replace', 'db.t', 'v1', '--snapshot-id', '999') + self.assertEqual(1, code) + self.assertIn("Failed to replace tag", err) + self.assertIn("Snapshot id '999' doesn't exist.", err) + + def test_replace_without_time_retained_drops_ttl(self): + self._run('tag', 'create', 'db.t', 'v1', '--time-retained', '12h') + _, _, code = self._run( + 'tag', 'replace', 'db.t', 'v1', '--snapshot-id', '1') + self.assertEqual(0, code) + got, _, _ = self._run('tag', 'get', 'db.t', 'v1') + self.assertIn("Snapshot ID: 1", got) + self.assertNotIn("Time Retained:", got) + # -- bad input ----------------------------------------------------------- def test_invalid_identifier(self): diff --git a/paimon-python/pypaimon/tests/rest/rest_tag_test.py b/paimon-python/pypaimon/tests/rest/rest_tag_test.py index 700bbc578c86..3fb2c438aeba 100644 --- a/paimon-python/pypaimon/tests/rest/rest_tag_test.py +++ b/paimon-python/pypaimon/tests/rest/rest_tag_test.py @@ -122,6 +122,25 @@ def test_delete_tag_not_exists(self): with self.assertRaises(TagNotExistException): self.rest_catalog.delete_tag(identifier, "absent") + def test_rename_tag_not_supported(self): + identifier = self._identifier() + self.rest_catalog.create_tag(identifier, "t1") + with self.assertRaises(NotImplementedError) as cm: + self.rest_catalog.rename_tag(identifier, "t1", "t2") + self.assertIn("rename_tag", str(cm.exception)) + # The server registry is unchanged; the old name still resolves. + response = self.rest_catalog.get_tag(identifier, "t1") + self.assertEqual(response.tag_name, "t1") + + def test_replace_tag_not_supported(self): + identifier = self._identifier() + self.rest_catalog.create_tag(identifier, "t1", snapshot_id=1) + with self.assertRaises(NotImplementedError) as cm: + self.rest_catalog.replace_tag(identifier, "t1", snapshot_id=1) + self.assertIn("replace_tag", str(cm.exception)) + response = self.rest_catalog.get_tag(identifier, "t1") + self.assertEqual(response.snapshot.id, 1) + # Note: the previous ``FilesystemCatalogTagInheritsNotImplementedTest`` class # has been removed because FileSystemCatalog now overrides the tag CRUD