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
9 changes: 7 additions & 2 deletions pyiceberg/transforms.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,12 +145,17 @@ def _cast_if_needed(arr: "ArrayLike") -> "ArrayLike":
else:
return arr

def _normalize_array(arr: "pa.Array") -> "pa.Array":
if pa.types.is_dictionary(arr.type):
return arr.dictionary_decode()
return arr

if isinstance(array, pa.Array):
return _cast_if_needed(transform_func(array, *args))
return _cast_if_needed(transform_func(_normalize_array(array), *args))
elif isinstance(array, pa.ChunkedArray):
result_chunks = []
for arr in array.iterchunks():
result_chunks.append(_cast_if_needed(transform_func(arr, *args)))
result_chunks.append(_cast_if_needed(transform_func(_normalize_array(arr), *args)))
return pa.chunked_array(result_chunks)
else:
raise ValueError(f"PyArrow array can only be of type pa.Array or pa.ChunkedArray, but found {type(array)}")
Expand Down
18 changes: 18 additions & 0 deletions tests/test_transforms.py
Original file line number Diff line number Diff line change
Expand Up @@ -1710,3 +1710,21 @@ def test_calling_pyarrow_transform_without_pyiceberg_core_installed_correctly_ra

with pytest.raises(NotInstalledError):
transform.pyarrow_transform(StringType())


def test_pyarrow_transforms_dictionary_encoded() -> None:
dict_arr = pa.DictionaryArray.from_arrays(pa.array([0, 1, 0, None]), pa.array(["foo", "bar"]))
raw_arr = pa.array(["foo", "bar", "foo", None])
bucket_transform = BucketTransform(num_buckets=10)
expected_bucket = bucket_transform.pyarrow_transform(StringType())(raw_arr)
assert bucket_transform.pyarrow_transform(StringType())(dict_arr) == expected_bucket
Comment on lines +1715 to +1720

chunked_dict = pa.chunked_array([dict_arr, dict_arr])
expected_chunked = pa.chunked_array([expected_bucket, expected_bucket])
assert bucket_transform.pyarrow_transform(StringType())(chunked_dict) == expected_chunked

truncate_transform = TruncateTransform(width=3)
dict_truncate_arr = pa.DictionaryArray.from_arrays(pa.array([0, 1, 0]), pa.array(["developer", "iceberg"]))
Comment on lines +1716 to +1727
raw_truncate_arr = pa.array(["developer", "iceberg", "developer"])
expected_truncate = truncate_transform.pyarrow_transform(StringType())(raw_truncate_arr)
assert truncate_transform.pyarrow_transform(StringType())(dict_truncate_arr) == expected_truncate