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
58 changes: 52 additions & 6 deletions python/pyarrow/interchange/from_dataframe.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@

from __future__ import annotations

import array
import sys
from typing import (
Any,
Tuple,
Expand Down Expand Up @@ -301,6 +303,7 @@ def parse_datetime_format_str(format_str):

def map_date_type(data_type):
"""Map column date type to pyarrow date type. """
# Buffer endianness conversion is handled in buffers_to_array.
kind, bit_width, f_string, _ = data_type

if kind == DtypeKind.DATETIME:
Expand Down Expand Up @@ -367,9 +370,12 @@ def buffers_to_array(
except TypeError:
offset_buff = None

# Construct a pyarrow Buffer
data_pa_buffer = pa.foreign_buffer(data_buff.ptr, data_buff.bufsize,
base=data_buff)
# Arrow buffers use native endianness. The interchange protocol allows
# producers to expose buffers in either byte order, so normalize fixed-
# width values before interpreting them as Arrow arrays.
data_pa_buffer = _buffer_with_native_endianness(
data_buff, data_type, allow_copy
)

# Construct a validity pyarrow Buffer, if applicable
if validity_buff:
Expand All @@ -394,9 +400,9 @@ def buffers_to_array(
_, offset_bit_width, _, _ = offset_dtype
# If an offset buffer exists, construct an offset pyarrow Buffer
# and add it to the construction of an array
offset_pa_buffer = pa.foreign_buffer(offset_buff.ptr,
offset_buff.bufsize,
base=offset_buff)
offset_pa_buffer = _buffer_with_native_endianness(
offset_buff, offset_dtype, allow_copy
)

if data_type[2] == 'U':
string_type = pa.large_string()
Expand All @@ -422,6 +428,46 @@ def buffers_to_array(
return array


def _buffer_with_native_endianness(buffer, dtype, allow_copy):
"""Wrap a protocol buffer, byte-swapping it when necessary."""
_, bit_width, _, endianness = dtype
native_endianness = "<" if sys.byteorder == "little" else ">"

if bit_width <= 8 or endianness in ("=", "|", native_endianness):
return pa.foreign_buffer(buffer.ptr, buffer.bufsize, base=buffer)

if endianness not in ("<", ">"):
raise ValueError(f"Unsupported endianness {endianness!r}")

if not allow_copy:
raise RuntimeError(
"Converting non-native endianness requires a copy which is "
"forbidden by allow_copy=False"
)

itemsize = bit_width // 8
if bit_width % 8 or buffer.bufsize % itemsize:
raise ValueError(
f"Buffer size {buffer.bufsize} is not a multiple of "
f"dtype width {bit_width}"
)

raw = pa.foreign_buffer(buffer.ptr, buffer.bufsize, base=buffer).to_pybytes()
typecode = next(
(code for code in "HIQ" if array.array(code).itemsize == itemsize),
None,
)
if typecode is None:
raise NotImplementedError(
f"No native array type with {itemsize}-byte elements is available"
)

native_values = array.array(typecode)
native_values.frombytes(raw)
native_values.byteswap()
return pa.py_buffer(native_values.tobytes())


def validity_buffer_from_mask(
validity_buff: BufferObject,
validity_dtype: Dtype,
Expand Down
153 changes: 152 additions & 1 deletion python/pyarrow/tests/interchange/test_conversion.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
# under the License.

from datetime import datetime as dt
import struct
import sys
import pyarrow as pa
import pytest

Expand All @@ -30,7 +32,8 @@
ColumnNullType,
DtypeKind,
)
from pyarrow.interchange.from_dataframe import _from_dataframe
from pyarrow.interchange.buffer import _PyArrowBuffer
from pyarrow.interchange.from_dataframe import _from_dataframe, buffers_to_array

try:
import pandas as pd
Expand Down Expand Up @@ -106,6 +109,154 @@ def test_offset_of_sliced_array():
# check_index=False, check_names=False)


def test_buffers_to_array_non_native_endian_numeric():
endianness = ">" if sys.byteorder == "little" else "<"
raw = struct.pack(f"{endianness}3i", 1, -2, 300)
dtype = (DtypeKind.INT, 32, "i", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(raw)), dtype),
"validity": None,
"offsets": None,
}

result = buffers_to_array(
buffers, dtype, 3, (ColumnNullType.NON_NULLABLE, None)
)
assert result.to_pylist() == [1, -2, 300]

with pytest.raises(RuntimeError, match="requires a copy"):
buffers_to_array(
buffers, dtype, 3, (ColumnNullType.NON_NULLABLE, None),
allow_copy=False,
)


def test_buffers_to_array_non_native_endian_sentinel_null():
endianness = ">" if sys.byteorder == "little" else "<"
raw = struct.pack(f"{endianness}3i", 1, -1, 300)
dtype = (DtypeKind.INT, 32, "i", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(raw)), dtype),
"validity": None,
"offsets": None,
}

result = buffers_to_array(
buffers, dtype, 3, (ColumnNullType.USE_SENTINEL, -1)
)
assert result.to_pylist() == [1, None, 300]


def test_buffers_to_array_non_native_endian_int16():
endianness = ">" if sys.byteorder == "little" else "<"
raw = struct.pack(f"{endianness}2h", 1, -2)
dtype = (DtypeKind.INT, 16, "s", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(raw)), dtype),
"validity": None,
"offsets": None,
}

result = buffers_to_array(
buffers, dtype, 2, (ColumnNullType.NON_NULLABLE, None)
)
assert result.to_pylist() == [1, -2]


def test_buffers_to_array_non_native_endian_int64():
endianness = ">" if sys.byteorder == "little" else "<"
raw = struct.pack(f"{endianness}2q", 1, -2)
dtype = (DtypeKind.INT, 64, "l", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(raw)), dtype),
"validity": None,
"offsets": None,
}

result = buffers_to_array(
buffers, dtype, 2, (ColumnNullType.NON_NULLABLE, None)
)
assert result.to_pylist() == [1, -2]


def test_buffers_to_array_int8_endianness_is_ignored():
endianness = ">" if sys.byteorder == "little" else "<"
raw = struct.pack(f"{endianness}2b", 1, -2)
dtype = (DtypeKind.INT, 8, "c", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(raw)), dtype),
"validity": None,
"offsets": None,
}

result = buffers_to_array(
buffers, dtype, 2, (ColumnNullType.NON_NULLABLE, None),
allow_copy=False,
)
assert result.to_pylist() == [1, -2]


@pytest.mark.pandas
def test_from_dataframe_non_native_endian_numeric():
df = pd.DataFrame({"value": np.array([1, 2, 300], dtype=">i4")})

table = pi.from_dataframe(df)

assert table["value"].to_pylist() == [1, 2, 300]


def test_buffers_to_array_unsupported_endianness():
dtype = (DtypeKind.INT, 32, "i", "?")
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(b"\x00" * 4)), dtype),
"validity": None,
"offsets": None,
}

with pytest.raises(ValueError, match="Unsupported endianness"):
buffers_to_array(
buffers, dtype, 1, (ColumnNullType.NON_NULLABLE, None)
)


def test_buffers_to_array_misaligned_buffer_size():
endianness = ">" if sys.byteorder == "little" else "<"
dtype = (DtypeKind.INT, 32, "i", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(b"\x00" * 3)), dtype),
"validity": None,
"offsets": None,
}

with pytest.raises(ValueError, match="not a multiple of dtype width"):
buffers_to_array(
buffers, dtype, 1, (ColumnNullType.NON_NULLABLE, None)
)


def test_buffers_to_array_non_native_endian_string_offsets():
endianness = ">" if sys.byteorder == "little" else "<"
offsets = struct.pack(f"{endianness}3i", 0, 1, 4)
data_dtype = (DtypeKind.STRING, 8, "u", "|")
offset_dtype = (DtypeKind.INT, 32, "i", endianness)
buffers = {
"data": (_PyArrowBuffer(pa.py_buffer(b"afoo")), data_dtype),
"validity": None,
"offsets": (_PyArrowBuffer(pa.py_buffer(offsets)), offset_dtype),
}

result = buffers_to_array(
buffers, data_dtype, 2, (ColumnNullType.NON_NULLABLE, None)
)
assert result.to_pylist() == ["a", "foo"]

with pytest.raises(RuntimeError, match="requires a copy"):
buffers_to_array(
buffers, data_dtype, 2,
(ColumnNullType.NON_NULLABLE, None), allow_copy=False,
)


@pytest.mark.pandas
@pytest.mark.parametrize(
"uint", [pa.uint8(), pa.uint16(), pa.uint32()]
Expand Down
Loading