From 4390268c011d5b188f999eb8faefa71e06f9528a Mon Sep 17 00:00:00 2001 From: Yuya Ebihara Date: Sat, 30 May 2026 09:50:02 +0900 Subject: [PATCH 1/3] Add support for environment context --- pyiceberg/environment_context.py | 43 +++++++++++++++++++ pyiceberg/table/snapshots.py | 4 ++ pyiceberg/view/metadata.py | 3 +- tests/integration/test_deletes.py | 3 ++ tests/integration/test_inspect_table.py | 5 +++ .../test_writes/test_partitioned_writes.py | 15 +++++++ tests/integration/test_writes/test_writes.py | 25 +++++++++++ tests/table/test_snapshots.py | 7 +++ tests/test_environment_context.py | 34 +++++++++++++++ 9 files changed, 138 insertions(+), 1 deletion(-) create mode 100644 pyiceberg/environment_context.py create mode 100644 tests/test_environment_context.py diff --git a/pyiceberg/environment_context.py b/pyiceberg/environment_context.py new file mode 100644 index 0000000000..8791ae247f --- /dev/null +++ b/pyiceberg/environment_context.py @@ -0,0 +1,43 @@ +# 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 importlib.metadata import version + + +class EnvironmentContext: + _PROPERTIES: dict[str, str] = { + "engine-name": "pyiceberg", + "engine-version": version("pyiceberg"), + } + + def __init__(self) -> None: + raise NotImplementedError("EnvironmentContext is a utility class and cannot be instantiated.") + + @classmethod + def get(cls) -> dict[str, str]: + """Return a read-only copy of all properties.""" + return cls._PROPERTIES.copy() + + @classmethod + def put(cls, key: str, value: str) -> None: + """Will add the given key/value pair in a global properties map.""" + cls._PROPERTIES[key] = value + + @classmethod + def remove(cls, key: str) -> str | None: + """Remove the key from the global properties map.""" + return cls._PROPERTIES.pop(key, None) diff --git a/pyiceberg/table/snapshots.py b/pyiceberg/table/snapshots.py index 5e9e519a01..cb67d427ca 100644 --- a/pyiceberg/table/snapshots.py +++ b/pyiceberg/table/snapshots.py @@ -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 @@ -409,6 +410,9 @@ def _update_totals(total_property: str, added_property: str, removed_property: s removed_property=REMOVED_EQUALITY_DELETES, ) + for key, value in EnvironmentContext.get().items(): + summary.__setitem__(key, value) + return summary diff --git a/pyiceberg/view/metadata.py b/pyiceberg/view/metadata.py index 33766040e3..648810b319 100644 --- a/pyiceberg/view/metadata.py +++ b/pyiceberg/view/metadata.py @@ -21,6 +21,7 @@ from pydantic import Field, RootModel, field_validator +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.schema import Schema from pyiceberg.typedef import IcebergBaseModel, Identifier, Properties from pyiceberg.typedef import ViewVersion as ViewVersionLiteral @@ -51,7 +52,7 @@ class ViewVersion(IcebergBaseModel): """ID of the schema for the view version""" timestamp_ms: int = Field(alias="timestamp-ms", default_factory=lambda: int(time.time() * 1000)) """Timestamp when the version was created (ms from epoch)""" - summary: dict[str, str] = Field(default_factory=dict) + summary: dict[str, str] = Field(default_factory=lambda: EnvironmentContext.get()) """A string to string map of summary metadata about the version""" representations: list[ViewRepresentation] = Field() """A list of representations for the view definition""" diff --git a/tests/integration/test_deletes.py b/tests/integration/test_deletes.py index 20205c59fb..68b8cff25e 100644 --- a/tests/integration/test_deletes.py +++ b/tests/integration/test_deletes.py @@ -23,6 +23,7 @@ from pyspark.sql import SparkSession from pyiceberg.catalog.rest import RestCatalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.expressions import AlwaysTrue, EqualTo, LessThanOrEqual from pyiceberg.manifest import ManifestContent, ManifestEntryStatus @@ -482,6 +483,8 @@ def test_partitioned_table_positional_deletes_sequence_number(spark: SparkSessio "total-files-size": snapshots[2].summary["total-files-size"], "total-position-deletes": "1", "total-records": "4", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), }, ) diff --git a/tests/integration/test_inspect_table.py b/tests/integration/test_inspect_table.py index 4d8dfbe9bb..637e5dd9b2 100644 --- a/tests/integration/test_inspect_table.py +++ b/tests/integration/test_inspect_table.py @@ -27,6 +27,7 @@ from pytest_lazy_fixtures import lf from pyiceberg.catalog import Catalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.expressions import ( And, @@ -277,6 +278,8 @@ def test_inspect_snapshots( ("total-files-size", str(file_size)), ("total-position-deletes", "0"), ("total-equality-deletes", "0"), + ("engine-name", "pyiceberg"), + ("engine-version", EnvironmentContext.get().get("engine-version")), ] # Delete @@ -290,6 +293,8 @@ def test_inspect_snapshots( ("total-files-size", "0"), ("total-position-deletes", "0"), ("total-equality-deletes", "0"), + ("engine-name", "pyiceberg"), + ("engine-version", EnvironmentContext.get().get("engine-version")), ] lhs = spark.table(f"{identifier}.snapshots").toPandas() diff --git a/tests/integration/test_writes/test_partitioned_writes.py b/tests/integration/test_writes/test_partitioned_writes.py index 1d1488255f..abfb9bad10 100644 --- a/tests/integration/test_writes/test_partitioned_writes.py +++ b/tests/integration/test_writes/test_partitioned_writes.py @@ -25,6 +25,7 @@ from pyspark.sql import SparkSession from pyiceberg.catalog import Catalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.partitioning import PartitionField, PartitionSpec from pyiceberg.schema import Schema @@ -498,6 +499,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": str(file_size), "total-position-deletes": "0", "total-records": "3", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[1] == { @@ -511,6 +514,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": str(file_size * 2), "total-position-deletes": "0", "total-records": "6", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[2] == { "removed-files-size": str(file_size * 2), @@ -523,6 +528,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": "0", "total-data-files": "0", "total-records": "0", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[3] == { "changed-partition-count": "3", @@ -535,6 +542,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": str(file_size), "total-data-files": "3", "total-records": "3", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[4] == { "changed-partition-count": "3", @@ -547,6 +556,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": str(file_size * 2), "total-data-files": "6", "total-records": "6", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert "removed-files-size" in summaries[5] assert "total-files-size" in summaries[5] @@ -561,6 +572,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": summaries[5]["total-files-size"], "total-data-files": "2", "total-records": "2", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert "added-files-size" in summaries[6] assert "total-files-size" in summaries[6] @@ -575,6 +588,8 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro "total-files-size": summaries[6]["total-files-size"], "total-data-files": "4", "total-records": "4", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } diff --git a/tests/integration/test_writes/test_writes.py b/tests/integration/test_writes/test_writes.py index 30fdd76ab7..dc4f45cae0 100644 --- a/tests/integration/test_writes/test_writes.py +++ b/tests/integration/test_writes/test_writes.py @@ -44,6 +44,7 @@ from pyiceberg.catalog import Catalog, load_catalog from pyiceberg.catalog.hive import HiveCatalog from pyiceberg.catalog.sql import SqlCatalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import CommitFailedException, NoSuchTableError from pyiceberg.expressions import And, EqualTo, GreaterThanOrEqual, In, LessThan, Not from pyiceberg.io.pyarrow import UnsupportedPyArrowTypeException, _dataframe_to_data_files @@ -231,6 +232,8 @@ def test_summaries(spark: SparkSession, session_catalog: Catalog, arrow_table_wi "total-files-size": str(file_size), "total-position-deletes": "0", "total-records": "3", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } # Append @@ -244,6 +247,8 @@ def test_summaries(spark: SparkSession, session_catalog: Catalog, arrow_table_wi "total-files-size": str(file_size * 2), "total-position-deletes": "0", "total-records": "6", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } # Delete @@ -257,6 +262,8 @@ def test_summaries(spark: SparkSession, session_catalog: Catalog, arrow_table_wi "total-files-size": "0", "total-position-deletes": "0", "total-records": "0", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } # Append @@ -270,6 +277,8 @@ def test_summaries(spark: SparkSession, session_catalog: Catalog, arrow_table_wi "total-files-size": str(file_size), "total-position-deletes": "0", "total-records": "3", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } @@ -326,6 +335,8 @@ def test_summaries_partial_overwrite(spark: SparkSession, session_catalog: Catal "total-files-size": summaries[0]["total-files-size"], "total-position-deletes": "0", "total-records": "5", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } # Java produces: # { @@ -367,6 +378,8 @@ def test_summaries_partial_overwrite(spark: SparkSession, session_catalog: Catal "total-files-size": summaries[1]["total-files-size"], "total-position-deletes": "0", "total-records": "4", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert len(tbl.scan().to_pandas()) == 4 @@ -831,6 +844,8 @@ def test_summaries_with_only_nulls( "total-files-size": "0", "total-position-deletes": "0", "total-records": "0", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[1] == { @@ -843,6 +858,8 @@ def test_summaries_with_only_nulls( "total-files-size": str(file_size), "total-position-deletes": "0", "total-records": "2", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[2] == { @@ -855,6 +872,8 @@ def test_summaries_with_only_nulls( "total-files-size": "0", "total-position-deletes": "0", "total-records": "0", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert summaries[3] == { @@ -864,6 +883,8 @@ def test_summaries_with_only_nulls( "total-files-size": "0", "total-position-deletes": "0", "total-records": "0", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } @@ -1156,6 +1177,8 @@ def test_inspect_snapshots( ("total-files-size", str(file_size)), ("total-position-deletes", "0"), ("total-equality-deletes", "0"), + ("engine-name", "pyiceberg"), + ("engine-version", EnvironmentContext.get().get("engine-version")), ] # Delete @@ -1169,6 +1192,8 @@ def test_inspect_snapshots( ("total-files-size", "0"), ("total-position-deletes", "0"), ("total-equality-deletes", "0"), + ("engine-name", "pyiceberg"), + ("engine-version", EnvironmentContext.get().get("engine-version")), ] lhs = spark.table(f"{identifier}.snapshots").toPandas() diff --git a/tests/table/test_snapshots.py b/tests/table/test_snapshots.py index 0f72b08087..2140da065e 100644 --- a/tests/table/test_snapshots.py +++ b/tests/table/test_snapshots.py @@ -25,6 +25,7 @@ import pytest from pyiceberg.catalog import Catalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import ValidationException from pyiceberg.io.pyarrow import _dataframe_to_data_files from pyiceberg.manifest import DataFile, DataFileContent, ManifestContent, ManifestFile @@ -326,6 +327,8 @@ def test_merge_snapshot_summaries_empty() -> None: "total-files-size": "0", "total-position-deletes": "0", "total-equality-deletes": "0", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), }, ) @@ -360,6 +363,8 @@ def test_merge_snapshot_summaries_new_summary() -> None: "total-files-size": "4", "total-position-deletes": "5", "total-equality-deletes": "3", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), }, ) @@ -402,6 +407,8 @@ def test_merge_snapshot_summaries_overwrite_summary() -> None: "total-files-size": "5", "total-position-deletes": "6", "total-equality-deletes": "4", + "engine-name": "pyiceberg", + "engine-version": EnvironmentContext.get().get("engine-version"), } assert actual.additional_properties == expected diff --git a/tests/test_environment_context.py b/tests/test_environment_context.py new file mode 100644 index 0000000000..fe0b44fd75 --- /dev/null +++ b/tests/test_environment_context.py @@ -0,0 +1,34 @@ +# 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 re + +from pyiceberg.environment_context import EnvironmentContext + + +def test_default_value() -> None: + actual = EnvironmentContext.get() + assert len(actual) == 2 + assert actual["engine-name"] == "pyiceberg" + assert re.match(r"^\d+\.\d+\.\d+", actual["engine-version"]) + + +def test_put_and_remove() -> None: + EnvironmentContext.put("test-key", "test-value") + assert EnvironmentContext.get()["test-key"] == "test-value" + + EnvironmentContext.remove("test-key") + assert "test-key" not in EnvironmentContext.get() From abaf725eb1c87122b47b65070e2e36634d86583a Mon Sep 17 00:00:00 2001 From: Yuya Ebihara Date: Thu, 25 Jun 2026 04:36:30 +0900 Subject: [PATCH 2/3] fixup! Add support for environment context --- pyiceberg/environment_context.py | 4 +- pyiceberg/table/snapshots.py | 2 +- pyiceberg/view/metadata.py | 3 +- tests/integration/test_deletes.py | 3 +- tests/integration/test_inspect_table.py | 54 +-- .../test_writes/test_partitioned_writes.py | 199 ++++++----- tests/integration/test_writes/test_writes.py | 313 +++++++++--------- tests/integration/test_writes/utils.py | 9 + tests/table/test_snapshots.py | 9 +- tests/test_environment_context.py | 28 +- 10 files changed, 317 insertions(+), 307 deletions(-) diff --git a/pyiceberg/environment_context.py b/pyiceberg/environment_context.py index 8791ae247f..e9da38c874 100644 --- a/pyiceberg/environment_context.py +++ b/pyiceberg/environment_context.py @@ -15,13 +15,13 @@ # specific language governing permissions and limitations # under the License. -from importlib.metadata import version +from pyiceberg import __version__ class EnvironmentContext: _PROPERTIES: dict[str, str] = { "engine-name": "pyiceberg", - "engine-version": version("pyiceberg"), + "engine-version": __version__, } def __init__(self) -> None: diff --git a/pyiceberg/table/snapshots.py b/pyiceberg/table/snapshots.py index cb67d427ca..60ffe8c577 100644 --- a/pyiceberg/table/snapshots.py +++ b/pyiceberg/table/snapshots.py @@ -411,7 +411,7 @@ def _update_totals(total_property: str, added_property: str, removed_property: s ) for key, value in EnvironmentContext.get().items(): - summary.__setitem__(key, value) + summary[key] = value return summary diff --git a/pyiceberg/view/metadata.py b/pyiceberg/view/metadata.py index 648810b319..33766040e3 100644 --- a/pyiceberg/view/metadata.py +++ b/pyiceberg/view/metadata.py @@ -21,7 +21,6 @@ from pydantic import Field, RootModel, field_validator -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.schema import Schema from pyiceberg.typedef import IcebergBaseModel, Identifier, Properties from pyiceberg.typedef import ViewVersion as ViewVersionLiteral @@ -52,7 +51,7 @@ class ViewVersion(IcebergBaseModel): """ID of the schema for the view version""" timestamp_ms: int = Field(alias="timestamp-ms", default_factory=lambda: int(time.time() * 1000)) """Timestamp when the version was created (ms from epoch)""" - summary: dict[str, str] = Field(default_factory=lambda: EnvironmentContext.get()) + summary: dict[str, str] = Field(default_factory=dict) """A string to string map of summary metadata about the version""" representations: list[ViewRepresentation] = Field() """A list of representations for the view definition""" diff --git a/tests/integration/test_deletes.py b/tests/integration/test_deletes.py index 68b8cff25e..e495cc6b4d 100644 --- a/tests/integration/test_deletes.py +++ b/tests/integration/test_deletes.py @@ -483,8 +483,7 @@ def test_partitioned_table_positional_deletes_sequence_number(spark: SparkSessio "total-files-size": snapshots[2].summary["total-files-size"], "total-position-deletes": "1", "total-records": "4", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), + **EnvironmentContext.get(), }, ) diff --git a/tests/integration/test_inspect_table.py b/tests/integration/test_inspect_table.py index 637e5dd9b2..9041aab921 100644 --- a/tests/integration/test_inspect_table.py +++ b/tests/integration/test_inspect_table.py @@ -27,7 +27,6 @@ from pytest_lazy_fixtures import lf from pyiceberg.catalog import Catalog -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.expressions import ( And, @@ -53,6 +52,7 @@ TimestampType, TimestamptzType, ) +from tests.integration.test_writes.utils import with_environment_context_tuples TABLE_SCHEMA = Schema( NestedField(field_id=1, name="bool", field_type=BooleanType(), required=False), @@ -268,34 +268,34 @@ def test_inspect_snapshots( assert file_size > 0 # Append - assert df["summary"][0].as_py() == [ - ("added-files-size", str(file_size)), - ("added-data-files", "1"), - ("added-records", "3"), - ("total-data-files", "1"), - ("total-delete-files", "0"), - ("total-records", "3"), - ("total-files-size", str(file_size)), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ("engine-name", "pyiceberg"), - ("engine-version", EnvironmentContext.get().get("engine-version")), - ] + assert df["summary"][0].as_py() == with_environment_context_tuples( + [ + ("added-files-size", str(file_size)), + ("added-data-files", "1"), + ("added-records", "3"), + ("total-data-files", "1"), + ("total-delete-files", "0"), + ("total-records", "3"), + ("total-files-size", str(file_size)), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] + ) # Delete - assert df["summary"][1].as_py() == [ - ("removed-files-size", str(file_size)), - ("deleted-data-files", "1"), - ("deleted-records", "3"), - ("total-data-files", "0"), - ("total-delete-files", "0"), - ("total-records", "0"), - ("total-files-size", "0"), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ("engine-name", "pyiceberg"), - ("engine-version", EnvironmentContext.get().get("engine-version")), - ] + assert df["summary"][1].as_py() == with_environment_context_tuples( + [ + ("removed-files-size", str(file_size)), + ("deleted-data-files", "1"), + ("deleted-records", "3"), + ("total-data-files", "0"), + ("total-delete-files", "0"), + ("total-records", "0"), + ("total-files-size", "0"), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] + ) lhs = spark.table(f"{identifier}.snapshots").toPandas() rhs = df.to_pandas() diff --git a/tests/integration/test_writes/test_partitioned_writes.py b/tests/integration/test_writes/test_partitioned_writes.py index abfb9bad10..2ca4d476dc 100644 --- a/tests/integration/test_writes/test_partitioned_writes.py +++ b/tests/integration/test_writes/test_partitioned_writes.py @@ -25,7 +25,6 @@ from pyspark.sql import SparkSession from pyiceberg.catalog import Catalog -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.partitioning import PartitionField, PartitionSpec from pyiceberg.schema import Schema @@ -43,7 +42,7 @@ from pyiceberg.types import ( StringType, ) -from utils import TABLE_SCHEMA, _create_table +from utils import TABLE_SCHEMA, _create_table, with_environment_context @pytest.mark.integration @@ -488,109 +487,109 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro file_size = int(summaries[0]["added-files-size"]) assert file_size > 0 - assert summaries[0] == { - "changed-partition-count": "3", - "added-data-files": "3", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "3", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "3", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[0] == with_environment_context( + { + "changed-partition-count": "3", + "added-data-files": "3", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "3", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "3", + } + ) - assert summaries[1] == { - "changed-partition-count": "3", - "added-data-files": "3", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "6", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size * 2), - "total-position-deletes": "0", - "total-records": "6", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } - assert summaries[2] == { - "removed-files-size": str(file_size * 2), - "changed-partition-count": "3", - "total-equality-deletes": "0", - "deleted-data-files": "6", - "total-position-deletes": "0", - "total-delete-files": "0", - "deleted-records": "6", - "total-files-size": "0", - "total-data-files": "0", - "total-records": "0", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } - assert summaries[3] == { - "changed-partition-count": "3", - "added-data-files": "3", - "total-equality-deletes": "0", - "added-records": "3", - "total-position-deletes": "0", - "added-files-size": str(file_size), - "total-delete-files": "0", - "total-files-size": str(file_size), - "total-data-files": "3", - "total-records": "3", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } - assert summaries[4] == { - "changed-partition-count": "3", - "added-data-files": "3", - "total-equality-deletes": "0", - "added-records": "3", - "total-position-deletes": "0", - "added-files-size": str(file_size), - "total-delete-files": "0", - "total-files-size": str(file_size * 2), - "total-data-files": "6", - "total-records": "6", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[1] == with_environment_context( + { + "changed-partition-count": "3", + "added-data-files": "3", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "6", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size * 2), + "total-position-deletes": "0", + "total-records": "6", + } + ) + assert summaries[2] == with_environment_context( + { + "removed-files-size": str(file_size * 2), + "changed-partition-count": "3", + "total-equality-deletes": "0", + "deleted-data-files": "6", + "total-position-deletes": "0", + "total-delete-files": "0", + "deleted-records": "6", + "total-files-size": "0", + "total-data-files": "0", + "total-records": "0", + } + ) + assert summaries[3] == with_environment_context( + { + "changed-partition-count": "3", + "added-data-files": "3", + "total-equality-deletes": "0", + "added-records": "3", + "total-position-deletes": "0", + "added-files-size": str(file_size), + "total-delete-files": "0", + "total-files-size": str(file_size), + "total-data-files": "3", + "total-records": "3", + } + ) + assert summaries[4] == with_environment_context( + { + "changed-partition-count": "3", + "added-data-files": "3", + "total-equality-deletes": "0", + "added-records": "3", + "total-position-deletes": "0", + "added-files-size": str(file_size), + "total-delete-files": "0", + "total-files-size": str(file_size * 2), + "total-data-files": "6", + "total-records": "6", + } + ) assert "removed-files-size" in summaries[5] assert "total-files-size" in summaries[5] - assert summaries[5] == { - "removed-files-size": summaries[5]["removed-files-size"], - "changed-partition-count": "2", - "total-equality-deletes": "0", - "deleted-data-files": "4", - "total-position-deletes": "0", - "total-delete-files": "0", - "deleted-records": "4", - "total-files-size": summaries[5]["total-files-size"], - "total-data-files": "2", - "total-records": "2", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[5] == with_environment_context( + { + "removed-files-size": summaries[5]["removed-files-size"], + "changed-partition-count": "2", + "total-equality-deletes": "0", + "deleted-data-files": "4", + "total-position-deletes": "0", + "total-delete-files": "0", + "deleted-records": "4", + "total-files-size": summaries[5]["total-files-size"], + "total-data-files": "2", + "total-records": "2", + } + ) assert "added-files-size" in summaries[6] assert "total-files-size" in summaries[6] - assert summaries[6] == { - "changed-partition-count": "2", - "added-data-files": "2", - "total-equality-deletes": "0", - "added-records": "2", - "total-position-deletes": "0", - "added-files-size": summaries[6]["added-files-size"], - "total-delete-files": "0", - "total-files-size": summaries[6]["total-files-size"], - "total-data-files": "4", - "total-records": "4", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[6] == with_environment_context( + { + "changed-partition-count": "2", + "added-data-files": "2", + "total-equality-deletes": "0", + "added-records": "2", + "total-position-deletes": "0", + "added-files-size": summaries[6]["added-files-size"], + "total-delete-files": "0", + "total-files-size": summaries[6]["total-files-size"], + "total-data-files": "4", + "total-records": "4", + } + ) @pytest.mark.integration diff --git a/tests/integration/test_writes/test_writes.py b/tests/integration/test_writes/test_writes.py index dc4f45cae0..03d546310f 100644 --- a/tests/integration/test_writes/test_writes.py +++ b/tests/integration/test_writes/test_writes.py @@ -44,7 +44,6 @@ from pyiceberg.catalog import Catalog, load_catalog from pyiceberg.catalog.hive import HiveCatalog from pyiceberg.catalog.sql import SqlCatalog -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import CommitFailedException, NoSuchTableError from pyiceberg.expressions import And, EqualTo, GreaterThanOrEqual, In, LessThan, Not from pyiceberg.io.pyarrow import UnsupportedPyArrowTypeException, _dataframe_to_data_files @@ -67,7 +66,7 @@ UUIDType, ) from pyiceberg.view.metadata import SQLViewRepresentation, ViewVersion -from utils import TABLE_SCHEMA, _create_table +from utils import TABLE_SCHEMA, _create_table, with_environment_context, with_environment_context_tuples @pytest.fixture(scope="session", autouse=True) @@ -222,64 +221,64 @@ def test_summaries(spark: SparkSession, session_catalog: Catalog, arrow_table_wi assert file_size > 0 # Append - assert summaries[0] == { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "1", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "3", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[0] == with_environment_context( + { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "1", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "3", + } + ) # Append - assert summaries[1] == { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "2", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size * 2), - "total-position-deletes": "0", - "total-records": "6", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[1] == with_environment_context( + { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "2", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size * 2), + "total-position-deletes": "0", + "total-records": "6", + } + ) # Delete - assert summaries[2] == { - "deleted-data-files": "2", - "deleted-records": "6", - "removed-files-size": str(file_size * 2), - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[2] == with_environment_context( + { + "deleted-data-files": "2", + "deleted-records": "6", + "removed-files-size": str(file_size * 2), + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } + ) # Append - assert summaries[3] == { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "1", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "3", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[3] == with_environment_context( + { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "1", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "3", + } + ) @pytest.mark.integration @@ -324,20 +323,20 @@ def test_summaries_partial_overwrite(spark: SparkSession, session_catalog: Catal # APPEND assert "added-files-size" in summaries[0] assert "total-files-size" in summaries[0] - assert summaries[0] == { - "added-data-files": "3", - "added-files-size": summaries[0]["added-files-size"], - "added-records": "5", - "changed-partition-count": "3", - "total-data-files": "3", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": summaries[0]["total-files-size"], - "total-position-deletes": "0", - "total-records": "5", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[0] == with_environment_context( + { + "added-data-files": "3", + "added-files-size": summaries[0]["added-files-size"], + "added-records": "5", + "changed-partition-count": "3", + "total-data-files": "3", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": summaries[0]["total-files-size"], + "total-position-deletes": "0", + "total-records": "5", + } + ) # Java produces: # { # "added-data-files": "1", @@ -364,23 +363,23 @@ def test_summaries_partial_overwrite(spark: SparkSession, session_catalog: Catal assert "added-files-size" in summaries[1] assert "removed-files-size" in summaries[1] assert "total-files-size" in summaries[1] - assert summaries[1] == { - "added-data-files": "1", - "added-files-size": summaries[1]["added-files-size"], - "added-records": "2", - "changed-partition-count": "1", - "deleted-data-files": "1", - "deleted-records": "3", - "removed-files-size": summaries[1]["removed-files-size"], - "total-data-files": "3", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": summaries[1]["total-files-size"], - "total-position-deletes": "0", - "total-records": "4", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[1] == with_environment_context( + { + "added-data-files": "1", + "added-files-size": summaries[1]["added-files-size"], + "added-records": "2", + "changed-partition-count": "1", + "deleted-data-files": "1", + "deleted-records": "3", + "removed-files-size": summaries[1]["removed-files-size"], + "total-data-files": "3", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": summaries[1]["total-files-size"], + "total-position-deletes": "0", + "total-records": "4", + } + ) assert len(tbl.scan().to_pandas()) == 4 @@ -837,55 +836,55 @@ def test_summaries_with_only_nulls( file_size = int(summaries[1]["added-files-size"]) assert file_size > 0 - assert summaries[0] == { - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[0] == with_environment_context( + { + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } + ) - assert summaries[1] == { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "2", - "total-data-files": "1", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "2", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[1] == with_environment_context( + { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "2", + "total-data-files": "1", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "2", + } + ) - assert summaries[2] == { - "deleted-data-files": "1", - "deleted-records": "2", - "removed-files-size": str(file_size), - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[2] == with_environment_context( + { + "deleted-data-files": "1", + "deleted-records": "2", + "removed-files-size": str(file_size), + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } + ) - assert summaries[3] == { - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), - } + assert summaries[3] == with_environment_context( + { + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } + ) @pytest.mark.integration @@ -1167,34 +1166,34 @@ def test_inspect_snapshots( assert file_size > 0 # Append - assert df["summary"][0].as_py() == [ - ("added-files-size", str(file_size)), - ("added-data-files", "1"), - ("added-records", "3"), - ("total-data-files", "1"), - ("total-delete-files", "0"), - ("total-records", "3"), - ("total-files-size", str(file_size)), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ("engine-name", "pyiceberg"), - ("engine-version", EnvironmentContext.get().get("engine-version")), - ] + assert df["summary"][0].as_py() == with_environment_context_tuples( + [ + ("added-files-size", str(file_size)), + ("added-data-files", "1"), + ("added-records", "3"), + ("total-data-files", "1"), + ("total-delete-files", "0"), + ("total-records", "3"), + ("total-files-size", str(file_size)), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] + ) # Delete - assert df["summary"][1].as_py() == [ - ("removed-files-size", str(file_size)), - ("deleted-data-files", "1"), - ("deleted-records", "3"), - ("total-data-files", "0"), - ("total-delete-files", "0"), - ("total-records", "0"), - ("total-files-size", "0"), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ("engine-name", "pyiceberg"), - ("engine-version", EnvironmentContext.get().get("engine-version")), - ] + assert df["summary"][1].as_py() == with_environment_context_tuples( + [ + ("removed-files-size", str(file_size)), + ("deleted-data-files", "1"), + ("deleted-records", "3"), + ("total-data-files", "0"), + ("total-delete-files", "0"), + ("total-records", "0"), + ("total-files-size", "0"), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] + ) lhs = spark.table(f"{identifier}.snapshots").toPandas() rhs = df.to_pandas() diff --git a/tests/integration/test_writes/utils.py b/tests/integration/test_writes/utils.py index 4ab54d97e7..967f37e957 100644 --- a/tests/integration/test_writes/utils.py +++ b/tests/integration/test_writes/utils.py @@ -20,6 +20,7 @@ import pyarrow as pa from pyiceberg.catalog import Catalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC, PartitionSpec from pyiceberg.schema import Schema @@ -79,3 +80,11 @@ def _create_table( tbl.append(d) return tbl + + +def with_environment_context(summary: dict[str, str]) -> dict[str, str]: + return {**summary, **EnvironmentContext.get()} + + +def with_environment_context_tuples(summary: list[tuple[str, str]]) -> list[tuple[str, str]]: + return summary + list(EnvironmentContext.get().items()) diff --git a/tests/table/test_snapshots.py b/tests/table/test_snapshots.py index 2140da065e..2bb5b4049f 100644 --- a/tests/table/test_snapshots.py +++ b/tests/table/test_snapshots.py @@ -327,8 +327,7 @@ def test_merge_snapshot_summaries_empty() -> None: "total-files-size": "0", "total-position-deletes": "0", "total-equality-deletes": "0", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), + **EnvironmentContext.get(), }, ) @@ -363,8 +362,7 @@ def test_merge_snapshot_summaries_new_summary() -> None: "total-files-size": "4", "total-position-deletes": "5", "total-equality-deletes": "3", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), + **EnvironmentContext.get(), }, ) @@ -407,8 +405,7 @@ def test_merge_snapshot_summaries_overwrite_summary() -> None: "total-files-size": "5", "total-position-deletes": "6", "total-equality-deletes": "4", - "engine-name": "pyiceberg", - "engine-version": EnvironmentContext.get().get("engine-version"), + **EnvironmentContext.get(), } assert actual.additional_properties == expected diff --git a/tests/test_environment_context.py b/tests/test_environment_context.py index fe0b44fd75..2ee3e6fb53 100644 --- a/tests/test_environment_context.py +++ b/tests/test_environment_context.py @@ -14,21 +14,29 @@ # KIND, either express or implied. See the License for the # specific language governing permissions and limitations # under the License. -import re - +from pyiceberg import __version__ from pyiceberg.environment_context import EnvironmentContext def test_default_value() -> None: - actual = EnvironmentContext.get() - assert len(actual) == 2 - assert actual["engine-name"] == "pyiceberg" - assert re.match(r"^\d+\.\d+\.\d+", actual["engine-version"]) + assert EnvironmentContext.get() == { + "engine-name": "pyiceberg", + "engine-version": __version__, + } -def test_put_and_remove() -> None: - EnvironmentContext.put("test-key", "test-value") - assert EnvironmentContext.get()["test-key"] == "test-value" +def test_get_returns_copy() -> None: + actual = EnvironmentContext.get() + actual["test-key"] = "test-value" - EnvironmentContext.remove("test-key") assert "test-key" not in EnvironmentContext.get() + + +def test_put_and_remove() -> None: + try: + EnvironmentContext.put("test-key", "test-value") + assert EnvironmentContext.get()["test-key"] == "test-value" + assert EnvironmentContext.remove("test-key") == "test-value" + assert "test-key" not in EnvironmentContext.get() + finally: + EnvironmentContext.remove("test-key") From 93ba3f19337151da506d2178f8463469e6acdc1d Mon Sep 17 00:00:00 2001 From: Kevin Liu Date: Sun, 13 Sep 2026 15:35:38 -0700 Subject: [PATCH 3/3] Keep environment context fixed and isolate legacy tests Return fresh engine metadata, copy only engine fields into new snapshot summaries, and cover writes and property precedence without changing existing summary assertions. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- mkdocs/docs/api.md | 3 + pyiceberg/environment_context.py | 29 +- pyiceberg/table/snapshots.py | 6 +- tests/conftest.py | 18 ++ tests/integration/test_deletes.py | 2 - tests/integration/test_inspect_table.py | 49 ++- .../test_writes/test_partitioned_writes.py | 184 ++++++----- tests/integration/test_writes/test_writes.py | 288 ++++++++---------- tests/integration/test_writes/utils.py | 9 - tests/table/test_snapshots.py | 46 ++- tests/test_environment_context.py | 28 +- 11 files changed, 323 insertions(+), 339 deletions(-) diff --git a/mkdocs/docs/api.md b/mkdocs/docs/api.md index 1e17e64f42..8afbca02da 100644 --- a/mkdocs/docs/api.md +++ b/mkdocs/docs/api.md @@ -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: diff --git a/pyiceberg/environment_context.py b/pyiceberg/environment_context.py index e9da38c874..0a4f4eccaf 100644 --- a/pyiceberg/environment_context.py +++ b/pyiceberg/environment_context.py @@ -19,25 +19,12 @@ class EnvironmentContext: - _PROPERTIES: dict[str, str] = { - "engine-name": "pyiceberg", - "engine-version": __version__, - } + """Environment context carrying the engine name and version for snapshot summaries.""" - def __init__(self) -> None: - raise NotImplementedError("EnvironmentContext is a utility class and cannot be instantiated.") - - @classmethod - def get(cls) -> dict[str, str]: - """Return a read-only copy of all properties.""" - return cls._PROPERTIES.copy() - - @classmethod - def put(cls, key: str, value: str) -> None: - """Will add the given key/value pair in a global properties map.""" - cls._PROPERTIES[key] = value - - @classmethod - def remove(cls, key: str) -> str | None: - """Remove the key from the global properties map.""" - return cls._PROPERTIES.pop(key, None) + @staticmethod + def get() -> dict[str, str]: + """Return a new dictionary containing only the engine name and version.""" + return { + "engine-name": "pyiceberg", + "engine-version": __version__, + } diff --git a/pyiceberg/table/snapshots.py b/pyiceberg/table/snapshots.py index 60ffe8c577..ee439c45ee 100644 --- a/pyiceberg/table/snapshots.py +++ b/pyiceberg/table/snapshots.py @@ -410,8 +410,10 @@ def _update_totals(total_property: str, added_property: str, removed_property: s removed_property=REMOVED_EQUALITY_DELETES, ) - for key, value in EnvironmentContext.get().items(): - summary[key] = value + 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"] return summary diff --git a/tests/conftest.py b/tests/conftest.py index cca9146c16..c070e81eac 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -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, @@ -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]: + """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. diff --git a/tests/integration/test_deletes.py b/tests/integration/test_deletes.py index e495cc6b4d..20205c59fb 100644 --- a/tests/integration/test_deletes.py +++ b/tests/integration/test_deletes.py @@ -23,7 +23,6 @@ from pyspark.sql import SparkSession from pyiceberg.catalog.rest import RestCatalog -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.expressions import AlwaysTrue, EqualTo, LessThanOrEqual from pyiceberg.manifest import ManifestContent, ManifestEntryStatus @@ -483,7 +482,6 @@ def test_partitioned_table_positional_deletes_sequence_number(spark: SparkSessio "total-files-size": snapshots[2].summary["total-files-size"], "total-position-deletes": "1", "total-records": "4", - **EnvironmentContext.get(), }, ) diff --git a/tests/integration/test_inspect_table.py b/tests/integration/test_inspect_table.py index 9041aab921..4d8dfbe9bb 100644 --- a/tests/integration/test_inspect_table.py +++ b/tests/integration/test_inspect_table.py @@ -52,7 +52,6 @@ TimestampType, TimestamptzType, ) -from tests.integration.test_writes.utils import with_environment_context_tuples TABLE_SCHEMA = Schema( NestedField(field_id=1, name="bool", field_type=BooleanType(), required=False), @@ -268,34 +267,30 @@ def test_inspect_snapshots( assert file_size > 0 # Append - assert df["summary"][0].as_py() == with_environment_context_tuples( - [ - ("added-files-size", str(file_size)), - ("added-data-files", "1"), - ("added-records", "3"), - ("total-data-files", "1"), - ("total-delete-files", "0"), - ("total-records", "3"), - ("total-files-size", str(file_size)), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ] - ) + assert df["summary"][0].as_py() == [ + ("added-files-size", str(file_size)), + ("added-data-files", "1"), + ("added-records", "3"), + ("total-data-files", "1"), + ("total-delete-files", "0"), + ("total-records", "3"), + ("total-files-size", str(file_size)), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] # Delete - assert df["summary"][1].as_py() == with_environment_context_tuples( - [ - ("removed-files-size", str(file_size)), - ("deleted-data-files", "1"), - ("deleted-records", "3"), - ("total-data-files", "0"), - ("total-delete-files", "0"), - ("total-records", "0"), - ("total-files-size", "0"), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ] - ) + assert df["summary"][1].as_py() == [ + ("removed-files-size", str(file_size)), + ("deleted-data-files", "1"), + ("deleted-records", "3"), + ("total-data-files", "0"), + ("total-delete-files", "0"), + ("total-records", "0"), + ("total-files-size", "0"), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] lhs = spark.table(f"{identifier}.snapshots").toPandas() rhs = df.to_pandas() diff --git a/tests/integration/test_writes/test_partitioned_writes.py b/tests/integration/test_writes/test_partitioned_writes.py index 2ca4d476dc..1d1488255f 100644 --- a/tests/integration/test_writes/test_partitioned_writes.py +++ b/tests/integration/test_writes/test_partitioned_writes.py @@ -42,7 +42,7 @@ from pyiceberg.types import ( StringType, ) -from utils import TABLE_SCHEMA, _create_table, with_environment_context +from utils import TABLE_SCHEMA, _create_table @pytest.mark.integration @@ -487,109 +487,95 @@ def test_summaries_with_null(spark: SparkSession, session_catalog: Catalog, arro file_size = int(summaries[0]["added-files-size"]) assert file_size > 0 - assert summaries[0] == with_environment_context( - { - "changed-partition-count": "3", - "added-data-files": "3", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "3", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "3", - } - ) + assert summaries[0] == { + "changed-partition-count": "3", + "added-data-files": "3", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "3", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "3", + } - assert summaries[1] == with_environment_context( - { - "changed-partition-count": "3", - "added-data-files": "3", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "6", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size * 2), - "total-position-deletes": "0", - "total-records": "6", - } - ) - assert summaries[2] == with_environment_context( - { - "removed-files-size": str(file_size * 2), - "changed-partition-count": "3", - "total-equality-deletes": "0", - "deleted-data-files": "6", - "total-position-deletes": "0", - "total-delete-files": "0", - "deleted-records": "6", - "total-files-size": "0", - "total-data-files": "0", - "total-records": "0", - } - ) - assert summaries[3] == with_environment_context( - { - "changed-partition-count": "3", - "added-data-files": "3", - "total-equality-deletes": "0", - "added-records": "3", - "total-position-deletes": "0", - "added-files-size": str(file_size), - "total-delete-files": "0", - "total-files-size": str(file_size), - "total-data-files": "3", - "total-records": "3", - } - ) - assert summaries[4] == with_environment_context( - { - "changed-partition-count": "3", - "added-data-files": "3", - "total-equality-deletes": "0", - "added-records": "3", - "total-position-deletes": "0", - "added-files-size": str(file_size), - "total-delete-files": "0", - "total-files-size": str(file_size * 2), - "total-data-files": "6", - "total-records": "6", - } - ) + assert summaries[1] == { + "changed-partition-count": "3", + "added-data-files": "3", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "6", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size * 2), + "total-position-deletes": "0", + "total-records": "6", + } + assert summaries[2] == { + "removed-files-size": str(file_size * 2), + "changed-partition-count": "3", + "total-equality-deletes": "0", + "deleted-data-files": "6", + "total-position-deletes": "0", + "total-delete-files": "0", + "deleted-records": "6", + "total-files-size": "0", + "total-data-files": "0", + "total-records": "0", + } + assert summaries[3] == { + "changed-partition-count": "3", + "added-data-files": "3", + "total-equality-deletes": "0", + "added-records": "3", + "total-position-deletes": "0", + "added-files-size": str(file_size), + "total-delete-files": "0", + "total-files-size": str(file_size), + "total-data-files": "3", + "total-records": "3", + } + assert summaries[4] == { + "changed-partition-count": "3", + "added-data-files": "3", + "total-equality-deletes": "0", + "added-records": "3", + "total-position-deletes": "0", + "added-files-size": str(file_size), + "total-delete-files": "0", + "total-files-size": str(file_size * 2), + "total-data-files": "6", + "total-records": "6", + } assert "removed-files-size" in summaries[5] assert "total-files-size" in summaries[5] - assert summaries[5] == with_environment_context( - { - "removed-files-size": summaries[5]["removed-files-size"], - "changed-partition-count": "2", - "total-equality-deletes": "0", - "deleted-data-files": "4", - "total-position-deletes": "0", - "total-delete-files": "0", - "deleted-records": "4", - "total-files-size": summaries[5]["total-files-size"], - "total-data-files": "2", - "total-records": "2", - } - ) + assert summaries[5] == { + "removed-files-size": summaries[5]["removed-files-size"], + "changed-partition-count": "2", + "total-equality-deletes": "0", + "deleted-data-files": "4", + "total-position-deletes": "0", + "total-delete-files": "0", + "deleted-records": "4", + "total-files-size": summaries[5]["total-files-size"], + "total-data-files": "2", + "total-records": "2", + } assert "added-files-size" in summaries[6] assert "total-files-size" in summaries[6] - assert summaries[6] == with_environment_context( - { - "changed-partition-count": "2", - "added-data-files": "2", - "total-equality-deletes": "0", - "added-records": "2", - "total-position-deletes": "0", - "added-files-size": summaries[6]["added-files-size"], - "total-delete-files": "0", - "total-files-size": summaries[6]["total-files-size"], - "total-data-files": "4", - "total-records": "4", - } - ) + assert summaries[6] == { + "changed-partition-count": "2", + "added-data-files": "2", + "total-equality-deletes": "0", + "added-records": "2", + "total-position-deletes": "0", + "added-files-size": summaries[6]["added-files-size"], + "total-delete-files": "0", + "total-files-size": summaries[6]["total-files-size"], + "total-data-files": "4", + "total-records": "4", + } @pytest.mark.integration diff --git a/tests/integration/test_writes/test_writes.py b/tests/integration/test_writes/test_writes.py index 03d546310f..30fdd76ab7 100644 --- a/tests/integration/test_writes/test_writes.py +++ b/tests/integration/test_writes/test_writes.py @@ -66,7 +66,7 @@ UUIDType, ) from pyiceberg.view.metadata import SQLViewRepresentation, ViewVersion -from utils import TABLE_SCHEMA, _create_table, with_environment_context, with_environment_context_tuples +from utils import TABLE_SCHEMA, _create_table @pytest.fixture(scope="session", autouse=True) @@ -221,64 +221,56 @@ def test_summaries(spark: SparkSession, session_catalog: Catalog, arrow_table_wi assert file_size > 0 # Append - assert summaries[0] == with_environment_context( - { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "1", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "3", - } - ) + assert summaries[0] == { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "1", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "3", + } # Append - assert summaries[1] == with_environment_context( - { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "2", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size * 2), - "total-position-deletes": "0", - "total-records": "6", - } - ) + assert summaries[1] == { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "2", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size * 2), + "total-position-deletes": "0", + "total-records": "6", + } # Delete - assert summaries[2] == with_environment_context( - { - "deleted-data-files": "2", - "deleted-records": "6", - "removed-files-size": str(file_size * 2), - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - } - ) + assert summaries[2] == { + "deleted-data-files": "2", + "deleted-records": "6", + "removed-files-size": str(file_size * 2), + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } # Append - assert summaries[3] == with_environment_context( - { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "3", - "total-data-files": "1", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "3", - } - ) + assert summaries[3] == { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "3", + "total-data-files": "1", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "3", + } @pytest.mark.integration @@ -323,20 +315,18 @@ def test_summaries_partial_overwrite(spark: SparkSession, session_catalog: Catal # APPEND assert "added-files-size" in summaries[0] assert "total-files-size" in summaries[0] - assert summaries[0] == with_environment_context( - { - "added-data-files": "3", - "added-files-size": summaries[0]["added-files-size"], - "added-records": "5", - "changed-partition-count": "3", - "total-data-files": "3", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": summaries[0]["total-files-size"], - "total-position-deletes": "0", - "total-records": "5", - } - ) + assert summaries[0] == { + "added-data-files": "3", + "added-files-size": summaries[0]["added-files-size"], + "added-records": "5", + "changed-partition-count": "3", + "total-data-files": "3", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": summaries[0]["total-files-size"], + "total-position-deletes": "0", + "total-records": "5", + } # Java produces: # { # "added-data-files": "1", @@ -363,23 +353,21 @@ def test_summaries_partial_overwrite(spark: SparkSession, session_catalog: Catal assert "added-files-size" in summaries[1] assert "removed-files-size" in summaries[1] assert "total-files-size" in summaries[1] - assert summaries[1] == with_environment_context( - { - "added-data-files": "1", - "added-files-size": summaries[1]["added-files-size"], - "added-records": "2", - "changed-partition-count": "1", - "deleted-data-files": "1", - "deleted-records": "3", - "removed-files-size": summaries[1]["removed-files-size"], - "total-data-files": "3", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": summaries[1]["total-files-size"], - "total-position-deletes": "0", - "total-records": "4", - } - ) + assert summaries[1] == { + "added-data-files": "1", + "added-files-size": summaries[1]["added-files-size"], + "added-records": "2", + "changed-partition-count": "1", + "deleted-data-files": "1", + "deleted-records": "3", + "removed-files-size": summaries[1]["removed-files-size"], + "total-data-files": "3", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": summaries[1]["total-files-size"], + "total-position-deletes": "0", + "total-records": "4", + } assert len(tbl.scan().to_pandas()) == 4 @@ -836,55 +824,47 @@ def test_summaries_with_only_nulls( file_size = int(summaries[1]["added-files-size"]) assert file_size > 0 - assert summaries[0] == with_environment_context( - { - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - } - ) + assert summaries[0] == { + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } - assert summaries[1] == with_environment_context( - { - "added-data-files": "1", - "added-files-size": str(file_size), - "added-records": "2", - "total-data-files": "1", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": str(file_size), - "total-position-deletes": "0", - "total-records": "2", - } - ) + assert summaries[1] == { + "added-data-files": "1", + "added-files-size": str(file_size), + "added-records": "2", + "total-data-files": "1", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": str(file_size), + "total-position-deletes": "0", + "total-records": "2", + } - assert summaries[2] == with_environment_context( - { - "deleted-data-files": "1", - "deleted-records": "2", - "removed-files-size": str(file_size), - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - } - ) + assert summaries[2] == { + "deleted-data-files": "1", + "deleted-records": "2", + "removed-files-size": str(file_size), + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } - assert summaries[3] == with_environment_context( - { - "total-data-files": "0", - "total-delete-files": "0", - "total-equality-deletes": "0", - "total-files-size": "0", - "total-position-deletes": "0", - "total-records": "0", - } - ) + assert summaries[3] == { + "total-data-files": "0", + "total-delete-files": "0", + "total-equality-deletes": "0", + "total-files-size": "0", + "total-position-deletes": "0", + "total-records": "0", + } @pytest.mark.integration @@ -1166,34 +1146,30 @@ def test_inspect_snapshots( assert file_size > 0 # Append - assert df["summary"][0].as_py() == with_environment_context_tuples( - [ - ("added-files-size", str(file_size)), - ("added-data-files", "1"), - ("added-records", "3"), - ("total-data-files", "1"), - ("total-delete-files", "0"), - ("total-records", "3"), - ("total-files-size", str(file_size)), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ] - ) + assert df["summary"][0].as_py() == [ + ("added-files-size", str(file_size)), + ("added-data-files", "1"), + ("added-records", "3"), + ("total-data-files", "1"), + ("total-delete-files", "0"), + ("total-records", "3"), + ("total-files-size", str(file_size)), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] # Delete - assert df["summary"][1].as_py() == with_environment_context_tuples( - [ - ("removed-files-size", str(file_size)), - ("deleted-data-files", "1"), - ("deleted-records", "3"), - ("total-data-files", "0"), - ("total-delete-files", "0"), - ("total-records", "0"), - ("total-files-size", "0"), - ("total-position-deletes", "0"), - ("total-equality-deletes", "0"), - ] - ) + assert df["summary"][1].as_py() == [ + ("removed-files-size", str(file_size)), + ("deleted-data-files", "1"), + ("deleted-records", "3"), + ("total-data-files", "0"), + ("total-delete-files", "0"), + ("total-records", "0"), + ("total-files-size", "0"), + ("total-position-deletes", "0"), + ("total-equality-deletes", "0"), + ] lhs = spark.table(f"{identifier}.snapshots").toPandas() rhs = df.to_pandas() diff --git a/tests/integration/test_writes/utils.py b/tests/integration/test_writes/utils.py index 967f37e957..4ab54d97e7 100644 --- a/tests/integration/test_writes/utils.py +++ b/tests/integration/test_writes/utils.py @@ -20,7 +20,6 @@ import pyarrow as pa from pyiceberg.catalog import Catalog -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import NoSuchTableError from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC, PartitionSpec from pyiceberg.schema import Schema @@ -80,11 +79,3 @@ def _create_table( tbl.append(d) return tbl - - -def with_environment_context(summary: dict[str, str]) -> dict[str, str]: - return {**summary, **EnvironmentContext.get()} - - -def with_environment_context_tuples(summary: list[tuple[str, str]]) -> list[tuple[str, str]]: - return summary + list(EnvironmentContext.get().items()) diff --git a/tests/table/test_snapshots.py b/tests/table/test_snapshots.py index 2bb5b4049f..294defdcca 100644 --- a/tests/table/test_snapshots.py +++ b/tests/table/test_snapshots.py @@ -24,8 +24,8 @@ import pyarrow as pa import pytest +from pyiceberg import __version__ from pyiceberg.catalog import Catalog -from pyiceberg.environment_context import EnvironmentContext from pyiceberg.exceptions import ValidationException from pyiceberg.io.pyarrow import _dataframe_to_data_files from pyiceberg.manifest import DataFile, DataFileContent, ManifestContent, ManifestFile @@ -327,7 +327,6 @@ def test_merge_snapshot_summaries_empty() -> None: "total-files-size": "0", "total-position-deletes": "0", "total-equality-deletes": "0", - **EnvironmentContext.get(), }, ) @@ -362,7 +361,6 @@ def test_merge_snapshot_summaries_new_summary() -> None: "total-files-size": "4", "total-position-deletes": "5", "total-equality-deletes": "3", - **EnvironmentContext.get(), }, ) @@ -405,7 +403,6 @@ def test_merge_snapshot_summaries_overwrite_summary() -> None: "total-files-size": "5", "total-position-deletes": "6", "total-equality-deletes": "4", - **EnvironmentContext.get(), } assert actual.additional_properties == expected @@ -671,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( diff --git a/tests/test_environment_context.py b/tests/test_environment_context.py index 2ee3e6fb53..09eea1c86d 100644 --- a/tests/test_environment_context.py +++ b/tests/test_environment_context.py @@ -18,25 +18,15 @@ from pyiceberg.environment_context import EnvironmentContext -def test_default_value() -> None: - assert EnvironmentContext.get() == { +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__, } - - -def test_get_returns_copy() -> None: - actual = EnvironmentContext.get() - actual["test-key"] = "test-value" - - assert "test-key" not in EnvironmentContext.get() - - -def test_put_and_remove() -> None: - try: - EnvironmentContext.put("test-key", "test-value") - assert EnvironmentContext.get()["test-key"] == "test-value" - assert EnvironmentContext.remove("test-key") == "test-value" - assert "test-key" not in EnvironmentContext.get() - finally: - EnvironmentContext.remove("test-key") + assert EnvironmentContext.get() == second