From 4b86925a67d91d7b8798586c6aa3c40d858e77b5 Mon Sep 17 00:00:00 2001 From: Minh Vu <38443830+fallintoplace@users.noreply.github.com> Date: Sun, 4 Oct 2026 00:49:00 +0200 Subject: [PATCH] Fix writes with time bucket partitions Preserve integer bucket values during time partition conversion to avoid an AttributeError. --- pyiceberg/partitioning.py | 6 ++++-- tests/io/test_pyarrow.py | 38 ++++++++++++++++++++++++++++++++++++-- tests/test_transforms.py | 18 +++++++++++++++++- 3 files changed, 57 insertions(+), 5 deletions(-) diff --git a/pyiceberg/partitioning.py b/pyiceberg/partitioning.py index 3d06287b6c..5654014fc1 100644 --- a/pyiceberg/partitioning.py +++ b/pyiceberg/partitioning.py @@ -546,8 +546,10 @@ def _(type: IcebergType, value: int | date | None) -> int | None: @_to_partition_representation.register(TimeType) -def _(type: IcebergType, value: time | None) -> int | None: - return time_to_micros(value) if value is not None else None +def _(type: IcebergType, value: int | time | None) -> int | None: + if value is None or isinstance(value, int): + return value + return time_to_micros(value) @_to_partition_representation.register(UUIDType) diff --git a/tests/io/test_pyarrow.py b/tests/io/test_pyarrow.py index 4d5d4431cb..0c47b2f3e5 100644 --- a/tests/io/test_pyarrow.py +++ b/tests/io/test_pyarrow.py @@ -22,7 +22,7 @@ import uuid import warnings from collections.abc import Iterator -from datetime import date, datetime, timezone +from datetime import date, datetime, time, timezone from pathlib import Path from typing import Any from unittest.mock import MagicMock, patch @@ -93,7 +93,7 @@ from pyiceberg.table import FileScanTask, TableProperties, WriteTask from pyiceberg.table.metadata import TableMetadataV2 from pyiceberg.table.name_mapping import create_mapping_from_schema -from pyiceberg.transforms import HourTransform, IdentityTransform +from pyiceberg.transforms import BucketTransform, HourTransform, IdentityTransform from pyiceberg.typedef import UTF8, Properties, Record, TableVersion from pyiceberg.types import ( BinaryType, @@ -2784,6 +2784,40 @@ def test_partition_for_demo() -> None: ) +@pytest.mark.parametrize("format_version", [1, 2]) +def test_append_time_bucket_partition(tmp_path: Path, format_version: int) -> None: + schema = Schema(NestedField(1, "event_time", TimeType(), required=False)) + spec = PartitionSpec(PartitionField(1, 1000, BucketTransform(16), "event_time_bucket")) + data = pa.table( + {"event_time": [time(0), time(12, 34, 56, 123456), time(23, 59, 59, 999999), None, time(0)]}, + schema=schema.as_arrow(), + ) + + with InMemoryCatalog("test", warehouse=tmp_path.as_uri()) as catalog: + catalog.create_namespace("default") + table = catalog.create_table( + "default.events", schema=schema, partition_spec=spec, properties={"format-version": str(format_version)} + ) + table.append(data) + + rows = table.scan().to_arrow()["event_time"].to_pylist() + assert rows.count(None) == 1 + assert sorted(value for value in rows if value is not None) == sorted( + value for value in data["event_time"].to_pylist() if value is not None + ) + + files = table.inspect.data_files().to_pylist() + assert {file["partition"]["event_time_bucket"]: file["record_count"] for file in files} == { + 12: 2, + 11: 1, + 8: 1, + None: 1, + } + for file in files: + bucket = file["partition"]["event_time_bucket"] + assert f"/event_time_bucket={bucket if bucket is not None else 'null'}/" in file["file_path"] + + def test_partition_for_nested_field() -> None: schema = Schema( NestedField(id=1, name="foo", field_type=StringType(), required=True), diff --git a/tests/test_transforms.py b/tests/test_transforms.py index 998d8a224d..634d8a64b7 100644 --- a/tests/test_transforms.py +++ b/tests/test_transforms.py @@ -17,7 +17,7 @@ # under the License. # pylint: disable=eval-used,protected-access,redefined-outer-name from collections.abc import Callable -from datetime import date, datetime +from datetime import date, datetime, time from decimal import Decimal from typing import Annotated, Any from uuid import UUID @@ -1643,6 +1643,22 @@ def test_to_partition_representation_timestamps(source_type: PrimitiveType, valu assert _to_partition_representation(source_type, value) == expected +@pytest.mark.parametrize( + "value, expected", + [ + pytest.param(None, None, id="none"), + pytest.param(time(0), 0, id="midnight"), + pytest.param(time(12, 34, 56, 123456), 45_296_123_456, id="microseconds"), + pytest.param(time(23, 59, 59, 999999), 86_399_999_999, id="end_of_day"), + pytest.param(0, 0, id="zero_int_passthrough"), + pytest.param(12, 12, id="bucket_int_passthrough"), + pytest.param(86_399_999_999, 86_399_999_999, id="time_int_passthrough"), + ], +) +def test_to_partition_representation_time(value: int | time | None, expected: int | None) -> None: + assert _to_partition_representation(TimeType(), value) == expected + + def test_to_partition_representation_unrecognized_type_raises() -> None: with pytest.raises(ValueError, match="Type not recognized"): _to_partition_representation(TimestampType(), "not-a-datetime-or-int")