Skip to content
Merged
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
3 changes: 3 additions & 0 deletions mkdocs/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -1405,6 +1405,9 @@ tbl.overwrite(df, snapshot_properties={"abc": "def"})
assert tbl.metadata.snapshots[-1].summary["abc"] == "def"
```

New snapshot summaries automatically include `engine-name` (`pyiceberg`) and `engine-version`
(the installed PyIceberg version). These values override same-named entries in `snapshot_properties`.

## Snapshot Management

Manage snapshots with operations through the `Table` API:
Expand Down
30 changes: 30 additions & 0 deletions pyiceberg/environment_context.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
# 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.

from pyiceberg import __version__


class EnvironmentContext:
"""Environment context carrying the engine name and version for snapshot summaries."""

@staticmethod
def get() -> dict[str, str]:
"""Return a new dictionary containing only the engine name and version."""
return {
"engine-name": "pyiceberg",
"engine-version": __version__,
}
6 changes: 6 additions & 0 deletions pyiceberg/table/snapshots.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@

from pydantic import Field, PrivateAttr, model_serializer

from pyiceberg.environment_context import EnvironmentContext
from pyiceberg.io import FileIO
from pyiceberg.manifest import DataFile, DataFileContent, ManifestFile, _manifests
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC, PartitionSpec
Expand Down Expand Up @@ -409,6 +410,11 @@ def _update_totals(total_property: str, added_property: str, removed_property: s
removed_property=REMOVED_EQUALITY_DELETES,
)

if context := EnvironmentContext.get():
# Defensively select only engine fields so future context additions cannot overwrite snapshot metadata.
summary["engine-name"] = context["engine-name"]
summary["engine-version"] = context["engine-version"]
Comment on lines +413 to +416

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

i think this is the best because we dont want to accidentally allow EnvironmentContext override other summary


return summary


Expand Down
18 changes: 18 additions & 0 deletions tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
from pytest_lazy_fixtures import lf

from pyiceberg.catalog import Catalog, load_catalog
from pyiceberg.environment_context import EnvironmentContext
from pyiceberg.expressions import BoundReference
from pyiceberg.io import (
ADLS_ACCOUNT_KEY,
Expand Down Expand Up @@ -104,12 +105,29 @@
from pyiceberg.io.pyarrow import PyArrowFileIO


_original_environment_context_get = EnvironmentContext.get


def pytest_collection_modifyitems(items: list[pytest.Item]) -> None:
for item in items:
if not any(item.iter_markers()):
item.add_marker("unmarked")


@pytest.fixture(autouse=True, scope="session")
def _disable_environment_context() -> Generator[None, None, None]:
Comment on lines +117 to +118

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this allows us to not make a lot of changes to existing tests

"""Disable engine metadata for existing tests, including session-scoped fixtures."""
with pytest.MonkeyPatch.context() as monkeypatch:
monkeypatch.setattr(EnvironmentContext, "get", staticmethod(dict))
yield


@pytest.fixture
def enable_environment_context(monkeypatch: pytest.MonkeyPatch) -> None:
"""Restore real engine metadata for tests that explicitly request it."""
monkeypatch.setattr(EnvironmentContext, "get", staticmethod(_original_environment_context_get))


@pytest.fixture(autouse=True, scope="session")
def _isolate_pyiceberg_config() -> None:
"""Make test runs ignore your local PyIceberg config.
Expand Down
42 changes: 42 additions & 0 deletions tests/table/test_snapshots.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import pyarrow as pa
import pytest

from pyiceberg import __version__
from pyiceberg.catalog import Catalog
from pyiceberg.exceptions import ValidationException
from pyiceberg.io.pyarrow import _dataframe_to_data_files
Expand Down Expand Up @@ -667,6 +668,47 @@ def overwrite_table(catalog: Catalog, arrow_table_simple: pa.Table) -> Table:
return table


def test_snapshot_writes_include_engine_metadata(
enable_environment_context: None, catalog: Catalog, arrow_table_simple: pa.Table
) -> None:
catalog.create_namespace("default")
table = catalog.create_table("default.engine_metadata", arrow_table_simple.schema)
table.append(arrow_table_simple)
table.overwrite(arrow_table_simple)
table.delete()

table.refresh()
snapshots = table.snapshots()
assert snapshots
for snapshot in snapshots:
assert snapshot.summary is not None
assert snapshot.summary["engine-name"] == "pyiceberg"
assert snapshot.summary["engine-version"] == __version__


def test_snapshot_engine_metadata_overrides_snapshot_properties(
enable_environment_context: None, catalog: Catalog, arrow_table_simple: pa.Table
) -> None:
catalog.create_namespace("default")
table = catalog.create_table("default.engine_metadata", arrow_table_simple.schema)
table.append(
arrow_table_simple,
snapshot_properties={
"engine-name": "custom-engine",
"engine-version": "custom-version",
"job-id": "snapshot-job",
},
)

table.refresh()
snapshot = table.current_snapshot()
assert snapshot is not None
assert snapshot.summary is not None
assert snapshot.summary["engine-name"] == "pyiceberg"
assert snapshot.summary["engine-version"] == __version__
assert snapshot.summary["job-id"] == "snapshot-job"


def _write_data_file(table: Table, rows: pa.Table) -> DataFile:
return next(
iter(
Expand Down
32 changes: 32 additions & 0 deletions tests/test_environment_context.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
# 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.
from pyiceberg import __version__
from pyiceberg.environment_context import EnvironmentContext


def test_get_returns_fresh_engine_metadata(enable_environment_context: None) -> None:
first = EnvironmentContext.get()
second = EnvironmentContext.get()
assert first is not second

first.clear()

assert second == {
"engine-name": "pyiceberg",
"engine-version": __version__,
}
assert EnvironmentContext.get() == second
Loading