Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/codeql.yml
Original file line number Diff line number Diff line change
Expand Up @@ -46,11 +46,11 @@ jobs:
persist-credentials: false

- name: Initialize CodeQL
uses: github/codeql-action/init@95e58e9a2cdfd71adc6e0353d5c52f41a045d225 # v4.35.2
uses: github/codeql-action/init@e46ed2cbd01164d986452f91f178727624ae40d7 # v4.35.3
with:
languages: actions

- name: Perform CodeQL Analysis
uses: github/codeql-action/analyze@95e58e9a2cdfd71adc6e0353d5c52f41a045d225 # v4.35.2
uses: github/codeql-action/analyze@e46ed2cbd01164d986452f91f178727624ae40d7 # v4.35.3
with:
category: "/language:actions"
5 changes: 4 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@
# under the License.
.PHONY: help install install-uv check-license lint \
test test-integration test-integration-setup test-integration-exec test-integration-cleanup test-integration-rebuild \
test-s3 test-adls test-gcs test-coverage coverage-report \
test-s3 test-adls test-gcs test-coverage coverage-report test test-notebook\
docs-serve docs-build notebook notebook-infra \
clean

Expand Down Expand Up @@ -150,6 +150,9 @@ coverage-report: ## Combine and report coverage
uv run $(PYTHON_ARG) coverage html
uv run $(PYTHON_ARG) coverage xml

test-notebook: ## Run notebook tests (pyiceberg_example and spark_integration_example) via papermill
$(TEST_RUNNER) pytest tests/notebooks/test_pyiceberg_example.py tests/notebooks/test_spark_integration_example.py -m notebook $(PYTEST_ARGS)

# ================
# Documentation
# ================
Expand Down
5 changes: 4 additions & 1 deletion dev/docker-compose-integration.yml
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,10 @@
services:
spark-iceberg:
image: pyiceberg-spark:latest
build: spark/
build:
context: spark/
args:
ICEBERG_MAVEN_MIRROR: https://repository.apache.org/content/repositories/orgapacheiceberg-1282
container_name: pyiceberg-spark
networks:
iceberg_net:
Expand Down
14 changes: 12 additions & 2 deletions dev/spark/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,12 @@ FROM apache/spark:${BASE_IMAGE_SPARK_VERSION}

# Dependency versions - keep these compatible
# Changing these will invalidate the JAR download cache layer
ARG ICEBERG_VERSION=1.10.1
ARG ICEBERG_VERSION=1.11.0
ARG ICEBERG_SPARK_RUNTIME_VERSION=4.0_2.13
ARG HADOOP_VERSION=3.4.1
ARG AWS_SDK_VERSION=2.24.6
ARG MAVEN_MIRROR=https://repo.maven.apache.org/maven2
ARG ICEBERG_MAVEN_MIRROR=${MAVEN_MIRROR}

USER root
WORKDIR ${SPARK_HOME}
Expand All @@ -38,11 +39,20 @@ RUN apt-get update -qq && \

# Download JARs with retry logic (most cacheable - only changes when versions change)
# This is the slowest step, so we do it before copying config files
# Iceberg JARs use ICEBERG_MAVEN_MIRROR (defaults to MAVEN_MIRROR, can be overridden for staging repos)
RUN set -e && \
cd "${SPARK_HOME}/jars" && \
for jar_path in \
"org/apache/iceberg/iceberg-spark-runtime-${ICEBERG_SPARK_RUNTIME_VERSION}/${ICEBERG_VERSION}/iceberg-spark-runtime-${ICEBERG_SPARK_RUNTIME_VERSION}-${ICEBERG_VERSION}.jar" \
"org/apache/iceberg/iceberg-aws-bundle/${ICEBERG_VERSION}/iceberg-aws-bundle-${ICEBERG_VERSION}.jar" \
"org/apache/iceberg/iceberg-aws-bundle/${ICEBERG_VERSION}/iceberg-aws-bundle-${ICEBERG_VERSION}.jar"; \
do \
jar_name=$(basename "${jar_path}") && \
curl -fsSL --retry 3 --retry-delay 5 \
-o "${jar_name}" \
"${ICEBERG_MAVEN_MIRROR}/${jar_path}" && \
chown spark:spark "${jar_name}"; \
done && \
for jar_path in \
"org/apache/hadoop/hadoop-aws/${HADOOP_VERSION}/hadoop-aws-${HADOOP_VERSION}.jar" \
"software/amazon/awssdk/bundle/${AWS_SDK_VERSION}/bundle-${AWS_SDK_VERSION}.jar"; \
do \
Expand Down
11 changes: 11 additions & 0 deletions mkdocs/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -1529,6 +1529,17 @@ catalog = load_catalog("default")
catalog.view_exists("default.bar")
```

## Register a view

To register a view using existing metadata:

```python
catalog.register_view(
identifier="docs_example.bids",
metadata_location="s3://warehouse/path/to/metadata.json"
)
```

## Table Statistics Management

Manage table statistics with operations through the `Table` API:
Expand Down
16 changes: 16 additions & 0 deletions pyiceberg/catalog/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -690,6 +690,22 @@ def update_namespace_properties(
ValueError: If removals and updates have overlapping keys.
"""

@abstractmethod
def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
"""Register a new view using existing metadata.

Args:
identifier (Union[str, Identifier]): View identifier for the view
metadata_location (str): The location to the metadata

Returns:
View: The newly registered view

Raises:
ViewAlreadyExistsError: If the view already exists.
TableAlreadyExistsError: If a table with the same name already exists.
"""

@abstractmethod
def drop_view(self, identifier: str | Identifier) -> None:
"""Drop a view.
Expand Down
3 changes: 3 additions & 0 deletions pyiceberg/catalog/bigquery_metastore.py
Original file line number Diff line number Diff line change
Expand Up @@ -305,6 +305,9 @@ def register_table(self, identifier: str | Identifier, metadata_location: str, o
def list_views(self, namespace: str | Identifier) -> list[Identifier]:
raise NotImplementedError

def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
raise NotImplementedError

def drop_view(self, identifier: str | Identifier) -> None:
raise NotImplementedError

Expand Down
3 changes: 3 additions & 0 deletions pyiceberg/catalog/dynamodb.py
Original file line number Diff line number Diff line change
Expand Up @@ -553,6 +553,9 @@ def create_view(
def list_views(self, namespace: str | Identifier) -> list[Identifier]:
raise NotImplementedError

def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
raise NotImplementedError

def drop_view(self, identifier: str | Identifier) -> None:
raise NotImplementedError

Expand Down
3 changes: 3 additions & 0 deletions pyiceberg/catalog/glue.py
Original file line number Diff line number Diff line change
Expand Up @@ -970,6 +970,9 @@ def create_view(
def list_views(self, namespace: str | Identifier) -> list[Identifier]:
raise NotImplementedError

def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
raise NotImplementedError

def drop_view(self, identifier: str | Identifier) -> None:
raise NotImplementedError

Expand Down
3 changes: 3 additions & 0 deletions pyiceberg/catalog/hive.py
Original file line number Diff line number Diff line change
Expand Up @@ -854,6 +854,9 @@ def update_namespace_properties(

return PropertiesUpdateSummary(removed=list(removed or []), updated=list(updated or []), missing=list(expected_to_change))

def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
raise NotImplementedError

def drop_view(self, identifier: str | Identifier) -> None:
raise NotImplementedError

Expand Down
3 changes: 3 additions & 0 deletions pyiceberg/catalog/noop.py
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,9 @@ def view_exists(self, identifier: str | Identifier) -> bool:
def namespace_exists(self, namespace: str | Identifier) -> bool:
raise NotImplementedError

def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
raise NotImplementedError

def drop_view(self, identifier: str | Identifier) -> None:
raise NotImplementedError

Expand Down
30 changes: 30 additions & 0 deletions pyiceberg/catalog/rest/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,7 @@ class Endpoints:
list_views: str = "namespaces/{namespace}/views"
load_view: str = "namespaces/{namespace}/views/{view}"
create_view: str = "namespaces/{namespace}/views"
register_view: str = "namespaces/{namespace}/register-view"
drop_view: str = "namespaces/{namespace}/views/{view}"
view_exists: str = "namespaces/{namespace}/views/{view}"
plan_table_scan: str = "namespaces/{namespace}/tables/{table}/plan"
Expand Down Expand Up @@ -183,6 +184,7 @@ class Capability:
V1_LIST_VIEWS = Endpoint(http_method=HttpMethod.GET, path=f"{API_PREFIX}/{Endpoints.list_views}")
V1_LOAD_VIEW = Endpoint(http_method=HttpMethod.GET, path=f"{API_PREFIX}/{Endpoints.load_view}")
V1_VIEW_EXISTS = Endpoint(http_method=HttpMethod.HEAD, path=f"{API_PREFIX}/{Endpoints.view_exists}")
V1_REGISTER_VIEW = Endpoint(http_method=HttpMethod.POST, path=f"{API_PREFIX}/{Endpoints.register_view}")
V1_DELETE_VIEW = Endpoint(http_method=HttpMethod.DELETE, path=f"{API_PREFIX}/{Endpoints.drop_view}")
V1_SUBMIT_TABLE_SCAN_PLAN = Endpoint(http_method=HttpMethod.POST, path=f"{API_PREFIX}/{Endpoints.plan_table_scan}")
V1_TABLE_SCAN_PLAN_TASKS = Endpoint(http_method=HttpMethod.POST, path=f"{API_PREFIX}/{Endpoints.fetch_scan_tasks}")
Expand Down Expand Up @@ -322,6 +324,11 @@ class RegisterTableRequest(IcebergBaseModel):
overwrite: bool


class RegisterViewRequest(IcebergBaseModel):
name: str
metadata_location: str = Field(..., alias="metadata-location")


class ConfigResponse(IcebergBaseModel):
defaults: Properties | None = Field(default_factory=dict)
overrides: Properties | None = Field(default_factory=dict)
Expand Down Expand Up @@ -1332,6 +1339,29 @@ def view_exists(self, identifier: str | Identifier) -> bool:

return False

@retry(**_RETRY_ARGS)
def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
self._check_endpoint(Capability.V1_REGISTER_VIEW)
namespace_and_view = self._split_identifier_for_path(identifier, IdentifierKind.VIEW)
namespace = namespace_and_view["namespace"]
view = namespace_and_view["view"]
if self.table_exists(identifier):
raise TableAlreadyExistsError(f"Table {namespace}.{view} already exists")

request = RegisterViewRequest(name=view, metadata_location=metadata_location)
serialized_json = request.model_dump_json().encode(UTF8)
response = self._session.post(
self.url(Endpoints.register_view, namespace=namespace),
data=serialized_json,
)
try:
response.raise_for_status()
except HTTPError as exc:
_handle_non_200_response(exc, {409: ViewAlreadyExistsError})

view_response = ViewResponse.model_validate_json(response.text)
return self._response_to_view(self.identifier_to_tuple(identifier), view_response)

@retry(**_RETRY_ARGS)
def drop_view(self, identifier: str) -> None:
self._check_endpoint(Capability.V1_DELETE_VIEW)
Expand Down
3 changes: 3 additions & 0 deletions pyiceberg/catalog/sql.py
Original file line number Diff line number Diff line change
Expand Up @@ -745,6 +745,9 @@ def list_views(self, namespace: str | Identifier) -> list[Identifier]:
def view_exists(self, identifier: str | Identifier) -> bool:
raise NotImplementedError

def register_view(self, identifier: str | Identifier, metadata_location: str) -> View:
raise NotImplementedError

def drop_view(self, identifier: str | Identifier) -> None:
raise NotImplementedError

Expand Down
2 changes: 2 additions & 0 deletions pyiceberg/table/inspect.py
Original file line number Diff line number Diff line change
Expand Up @@ -404,6 +404,7 @@ def _get_all_manifests_schema(self) -> pa.Schema:

all_manifests_schema = self._get_manifests_schema()
all_manifests_schema = all_manifests_schema.append(pa.field("reference_snapshot_id", pa.int64(), nullable=False))
all_manifests_schema = all_manifests_schema.append(pa.field("key_metadata", pa.binary(), nullable=True))
return all_manifests_schema

def _generate_manifests_table(self, snapshot: Snapshot | None, is_all_manifests_table: bool = False) -> pa.Table:
Expand Down Expand Up @@ -468,6 +469,7 @@ def _partition_summaries_to_rows(
}
if is_all_manifests_table:
manifest_row["reference_snapshot_id"] = snapshot.snapshot_id
manifest_row["key_metadata"] = manifest.key_metadata
manifests.append(manifest_row)

return pa.Table.from_pylist(
Expand Down
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,9 @@ dev = [
"google-cloud-bigquery>=3.33.0,<4",
"pyarrow-stubs>=20.0.0.20251107", # Remove when pyarrow >= 23.0.0 https://github.com/apache/arrow/pull/47609
"sqlalchemy>=2.0.18,<3",
"papermill>=2.6.0",
"nbformat>=5.10.0",
"ipykernel>=6.29.0",
]
# for mkdocs
docs = [
Expand Down Expand Up @@ -161,6 +164,7 @@ markers = [
"integration: marks integration tests against Apache Spark",
"gcs: marks a test as requiring access to gcs compliant storage (use with --gs.token, --gs.project, and --gs.endpoint)",
"benchmark: collection of tests to validate read/write performance before and after a change",
"notebook: marks tests that execute Jupyter notebooks via papermill",
]

# Turns a warning into an error
Expand Down
67 changes: 67 additions & 0 deletions tests/catalog/test_rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,7 @@
Capability.V1_LIST_VIEWS,
Capability.V1_LOAD_VIEW,
Capability.V1_VIEW_EXISTS,
Capability.V1_REGISTER_VIEW,
Capability.V1_DELETE_VIEW,
Capability.V1_SUBMIT_TABLE_SCAN_PLAN,
Capability.V1_TABLE_SCAN_PLAN_TASKS,
Expand Down Expand Up @@ -2182,6 +2183,72 @@ def test_table_identifier_in_commit_table_request(
)


def test_register_view_200(rest_mock: Mocker, example_view_metadata_rest_json: dict[str, Any]) -> None:
rest_mock.head(
f"{TEST_URI}v1/namespaces/default/tables/registered_view",
status_code=404,
request_headers=TEST_HEADERS,
)
rest_mock.post(
f"{TEST_URI}v1/namespaces/default/register-view",
json=example_view_metadata_rest_json,
status_code=200,
request_headers=TEST_HEADERS,
)

catalog = RestCatalog("rest", uri=TEST_URI, token=TEST_TOKEN)
actual = catalog.register_view(
identifier=("default", "registered_view"), metadata_location="s3://warehouse/database/view/metadata.json"
)
expected = View(
identifier=("default", "registered_view"),
metadata=ViewMetadata(**example_view_metadata_rest_json["metadata"]),
)
assert actual == expected


def test_register_view_409_view(rest_mock: Mocker) -> None:
rest_mock.head(
f"{TEST_URI}v1/namespaces/default/tables/registered_view",
status_code=404,
request_headers=TEST_HEADERS,
)
rest_mock.post(
f"{TEST_URI}v1/namespaces/default/register-view",
json={
"error": {
"message": "View already exists: default.view in warehouse 8bcb0838-50fc-472d-9ddb-8feb89ef5f1e",
"type": "AlreadyExistsException",
"code": 409,
}
},
status_code=409,
request_headers=TEST_HEADERS,
)

catalog = RestCatalog("rest", uri=TEST_URI, token=TEST_TOKEN)
with pytest.raises(ViewAlreadyExistsError) as e:
catalog.register_view(
identifier=("default", "registered_view"), metadata_location="s3://warehouse/database/view/metadata.json"
)
assert "View already exists" in str(e.value)


def test_register_view_409_table(rest_mock: Mocker) -> None:
rest_mock.head(
f"{TEST_URI}v1/namespaces/default/tables/registered_view",
status_code=200,
request_headers=TEST_HEADERS,
)

catalog = RestCatalog("rest", uri=TEST_URI, token=TEST_TOKEN)
with pytest.raises(TableAlreadyExistsError) as e:
catalog.register_view(
identifier=("default", "registered_view"), metadata_location="s3://warehouse/database/view/metadata.json"
)
assert "Table default.registered_view already exists" in str(e.value)


def test_drop_view_invalid_namespace(rest_mock: Mocker) -> None:
view = "view"
with pytest.raises(NoSuchIdentifierError) as e:
Expand Down
1 change: 1 addition & 0 deletions tests/integration/test_inspect_table.py
Original file line number Diff line number Diff line change
Expand Up @@ -1012,6 +1012,7 @@ def test_inspect_all_manifests(spark: SparkSession, session_catalog: Catalog, fo
"deleted_delete_files_count",
"partition_summaries",
"reference_snapshot_id",
"key_metadata",
]

int_cols = [
Expand Down
2 changes: 1 addition & 1 deletion tests/integration/test_writes/test_writes.py
Original file line number Diff line number Diff line change
Expand Up @@ -2027,7 +2027,7 @@ def test_write_optional_list(session_catalog: Catalog) -> None:
required=False,
),
)
session_catalog.create_table_if_not_exists(identifier, schema)
_create_table(session_catalog, identifier, schema=schema)

df_1 = pa.Table.from_pylist(
[
Expand Down
Loading
Loading