diff --git a/pyiceberg/catalog/__init__.py b/pyiceberg/catalog/__init__.py index 95ceaa539f..bba3244e72 100644 --- a/pyiceberg/catalog/__init__.py +++ b/pyiceberg/catalog/__init__.py @@ -96,6 +96,7 @@ LOCATION = "location" EXTERNAL_TABLE = "EXTERNAL_TABLE" BOTOCORE_SESSION = "botocore_session" +ICEBERG_VIEW = "iceberg-view" TABLE_METADATA_FILE_NAME_REGEX = re.compile( r""" diff --git a/pyiceberg/catalog/hive.py b/pyiceberg/catalog/hive.py index 181f9d4661..d2787f5b1f 100644 --- a/pyiceberg/catalog/hive.py +++ b/pyiceberg/catalog/hive.py @@ -55,6 +55,7 @@ from pyiceberg.catalog import ( EXTERNAL_TABLE, ICEBERG, + ICEBERG_VIEW, LOCATION, METADATA_LOCATION, TABLE_TYPE, @@ -480,7 +481,29 @@ def register_table(self, identifier: str | Identifier, metadata_location: str, o @override def list_views(self, namespace: str | Identifier) -> list[Identifier]: - raise NotImplementedError + """List Iceberg views under the given namespace in the catalog. + + Args: + namespace: Database to list. + + Returns: + List[Identifier]: list of views identifiers. + + Raises: + NoSuchNamespaceError: If a namespace with the given name does not exist, or the identifier is invalid. + """ + database_name = self.identifier_to_database(namespace, NoSuchNamespaceError) + try: + with self._client as open_client: + return [ + (database_name, table.tableName) + for table in open_client.get_table_objects_by_name( + dbname=database_name, tbl_names=open_client.get_all_tables(db_name=database_name) + ) + if table.parameters.get(TABLE_TYPE, "").lower() == ICEBERG_VIEW + ] + except (MetaException, NoSuchObjectException) as e: + raise NoSuchNamespaceError(f"Database does not exists: {database_name}") from e @override def view_exists(self, identifier: str | Identifier) -> bool: diff --git a/tests/catalog/test_hive.py b/tests/catalog/test_hive.py index 09bb5ab920..84c2bf6572 100644 --- a/tests/catalog/test_hive.py +++ b/tests/catalog/test_hive.py @@ -1144,6 +1144,50 @@ def test_create_database_already_exists() -> None: assert "Database default already exists" in str(exc_info.value) +def test_list_views(hive_table: HiveTable) -> None: + catalog = HiveCatalog(HIVE_CATALOG_NAME, uri=HIVE_METASTORE_FAKE_URL) + + tbl1 = deepcopy(hive_table) + tbl1.tableName = "table1" + tbl1.dbName = "database" + tbl1.parameters["table_type"] = "iceberg-view" + tbl2 = deepcopy(hive_table) + tbl2.tableName = "table2" + tbl2.dbName = "database" + tbl2.parameters["table_type"] = "iceberg-view" + tbl3 = deepcopy(hive_table) + tbl3.tableName = "table3" + tbl3.dbName = "database" + tbl3.parameters["table_type"] = "iceberg" + tbl4 = deepcopy(hive_table) + tbl4.tableName = "table4" + tbl4.dbName = "database" + tbl4.parameters.pop("table_type") + + catalog._client = MagicMock() + catalog._client.__enter__().get_all_tables.return_value = ["table1", "table2", "table3", "table4"] + catalog._client.__enter__().get_table_objects_by_name.return_value = [tbl1, tbl2, tbl3, tbl4] + + got_tables = catalog.list_views("database") + assert got_tables == [("database", "table1"), ("database", "table2")] + catalog._client.__enter__().get_all_tables.assert_called_with(db_name="database") + catalog._client.__enter__().get_table_objects_by_name.assert_called_with( + dbname="database", tbl_names=["table1", "table2", "table3", "table4"] + ) + + +def test_list_views_to_namespace_does_not_exists(hive_table: HiveTable) -> None: + catalog = HiveCatalog(HIVE_CATALOG_NAME, uri=HIVE_METASTORE_FAKE_URL) + + catalog._client = MagicMock() + catalog._client.__enter__().get_all_tables.side_effect = NoSuchObjectException(message="does_not_exists") + + with pytest.raises(NoSuchNamespaceError) as exc_info: + catalog.list_views("database") + + assert "Database does not exists: database" in str(exc_info.value) + + def test_load_namespace_properties(hive_database: HiveDatabase) -> None: catalog = HiveCatalog(HIVE_CATALOG_NAME, uri=HIVE_METASTORE_FAKE_URL) diff --git a/tests/integration/test_hive_catalog.py b/tests/integration/test_hive_catalog.py new file mode 100644 index 0000000000..cd0e56f422 --- /dev/null +++ b/tests/integration/test_hive_catalog.py @@ -0,0 +1,76 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +import time + +import pytest +from pyspark.sql import SparkSession + +from pyiceberg.catalog.hive import HiveCatalog + + +@pytest.mark.integration +def test_list_views( + session_catalog_hive: HiveCatalog, + spark: SparkSession, +) -> None: + """ + Verify that a view created by Spark through the Iceberg Hive catalog + can be discovered by PyIceberg HiveCatalog.list_views(). + + The test also verifies that: + - Iceberg tables are not returned as views. + - Multiple views are returned. + - Returned identifiers use the expected (namespace, view_name) format. + """ + suffix = int(time.time()) + catalog_name = "hive" + namespace = "default" + table_name = f"table_{suffix}" + first_view_name = f"first_view_{suffix}" + second_view_name = f"second_view_{suffix}" + table_identifier = f"{catalog_name}.{namespace}.{table_name}" + first_view_identifier = f"{catalog_name}.{namespace}.{first_view_name}" + second_view_identifier = f"{catalog_name}.{namespace}.{second_view_name}" + + spark.sql(f""" + CREATE TABLE {table_identifier} ( + id INTEGER, + name STRING, + dt DATE + ) + USING iceberg + """) + + spark.sql(f""" + CREATE VIEW {first_view_identifier} AS + SELECT id, name + FROM {table_identifier} + """) + + spark.sql(f""" + CREATE VIEW {second_view_identifier} AS + SELECT id, name, dt + FROM {table_identifier} + """) + + views = set(session_catalog_hive.list_views(namespace)) + + assert (namespace, first_view_name) in views + assert (namespace, second_view_name) in views + + # A table in the same namespace must not be returned as a view. + assert (namespace, table_name) not in views