diff --git a/vortex-python/python/vortex/_lib/arrays.pyi b/vortex-python/python/vortex/_lib/arrays.pyi index a0c8317e9db..1c1d74bd093 100644 --- a/vortex-python/python/vortex/_lib/arrays.pyi +++ b/vortex-python/python/vortex/_lib/arrays.pyi @@ -25,7 +25,9 @@ class Array: ) -> Array: ... @staticmethod def from_range(obj: range, *, dtype: DType | None = None) -> Array: ... - def to_arrow_array(self, *, arrow_type: pa.DataType | None = None) -> pa.Array[pa.Scalar[pa.DataType]]: ... + def to_arrow_array( + self, *, arrow_type: pa.DataType | None = None, combine_chunks: bool = False + ) -> pa.Array[pa.Scalar[pa.DataType]]: ... def __vortex_array_metadata__(self) -> tuple[str, bytes, int, bytes, list[object], list[object]]: ... @property def id(self) -> str: ... @@ -38,7 +40,7 @@ class Array: def take(self, indices: Array) -> Array: ... def slice(self, start: int, end: int) -> Array: ... def display_tree(self) -> str: ... - def to_arrow_table(self) -> pa.Table: ... + def to_arrow_table(self, *, combine_chunks: bool = False) -> pa.Table: ... def to_numpy(self, *, zero_copy_only: bool = True) -> np.ndarray: ... def to_pandas(self) -> pd.DataFrame: ... def to_polars_dataframe(self) -> pl.DataFrame: ... diff --git a/vortex-python/python/vortex/arrays.py b/vortex-python/python/vortex/arrays.py index 200df5cb1f2..165879e48c8 100644 --- a/vortex-python/python/vortex/arrays.py +++ b/vortex-python/python/vortex/arrays.py @@ -62,7 +62,7 @@ def arrow_table_from_struct_array( return pyarrow.Table.from_struct_array(array) -def _Array_to_arrow_table(self: _arrays.Array) -> pyarrow.Table: +def _Array_to_arrow_table(self: _arrays.Array, *, combine_chunks: bool = False) -> pyarrow.Table: """Construct an Arrow table from this Vortex array. .. seealso:: @@ -73,6 +73,12 @@ def _Array_to_arrow_table(self: _arrays.Array) -> pyarrow.Table: Only struct-typed arrays can be converted to Arrow tables. + Parameters + ---------- + combine_chunks : :class:`bool`, optional + If ``True``, a chunked Vortex array is concatenated natively during export so that every + column of the resulting table has a single chunk. Defaults to ``False``. + Returns ------- @@ -96,7 +102,7 @@ def _Array_to_arrow_table(self: _arrays.Array) -> pyarrow.Table: age: [[25,31,33,57]] """ - array = self.to_arrow_array() + array = self.to_arrow_array(combine_chunks=combine_chunks) assert isinstance(array, pyarrow.StructArray | pyarrow.ChunkedArray) return arrow_table_from_struct_array(array) diff --git a/vortex-python/src/arrays/mod.rs b/vortex-python/src/arrays/mod.rs index 06e94f8d464..d6d287cf1f9 100644 --- a/vortex-python/src/arrays/mod.rs +++ b/vortex-python/src/arrays/mod.rs @@ -436,10 +436,17 @@ impl PyArray { /// arrow_type : :class:`pyarrow.DataType`, optional /// The Arrow type to return. By default, UTF-8 data returns a ``StringViewArray`` and /// binary data returns a ``BinaryViewArray``. + /// combine_chunks : :class:`bool`, optional + /// If ``True``, a chunked Vortex array is exported as a single contiguous + /// :class:`pyarrow.Array` instead of a :class:`pyarrow.ChunkedArray`. The chunks are + /// concatenated natively during export, which is faster than calling + /// :meth:`pyarrow.ChunkedArray.combine_chunks` on the result. Defaults to ``False``. /// /// Returns /// ------- - /// :class:`pyarrow.Array` + /// :class:`pyarrow.Array` or :class:`pyarrow.ChunkedArray` + /// A :class:`pyarrow.ChunkedArray` is returned only for chunked arrays when + /// ``combine_chunks`` is ``False``. /// /// Examples /// -------- @@ -468,10 +475,26 @@ impl PyArray { /// "world" /// ] /// ``` - #[pyo3(signature = (*, arrow_type = None))] + /// + /// Export a chunked array as a single contiguous Arrow array: + /// + /// ```python + /// >>> import pyarrow + /// >>> import vortex as vx + /// >>> chunked = vx.array(pyarrow.chunked_array([[1, 2], [3]])) + /// >>> chunked.to_arrow_array(combine_chunks=True) + /// + /// [ + /// 1, + /// 2, + /// 3 + /// ] + /// ``` + #[pyo3(signature = (*, arrow_type = None, combine_chunks = false))] fn to_arrow_array<'py>( self_: &'py Bound<'py, Self>, arrow_type: Option<&Bound<'py, PyAny>>, + combine_chunks: bool, ) -> PyVortexResult> { // NOTE(ngates): for struct arrays, we could also return a RecordBatchStreamReader. let array = PyArrayRef::extract(self_.as_any().as_borrowed())?.into_inner(); @@ -481,7 +504,9 @@ impl PyArray { .transpose()? .map(|data_type| Field::new("", data_type, array.dtype().is_nullable())); - if let Some(chunked_array) = array.as_opt::() { + if let Some(chunked_array) = array.as_opt::() + && !combine_chunks + { // We figure out a single Arrow Data Type to convert all chunks into, otherwise // the preferred type of each chunk may be different. let inferred_field; diff --git a/vortex-python/test/test_array.py b/vortex-python/test/test_array.py index 9d20adba5d6..a10edd14938 100644 --- a/vortex-python/test/test_array.py +++ b/vortex-python/test/test_array.py @@ -85,3 +85,31 @@ def test_unsupported_arrow_type_raises_value_error(arrow_type: pa.DataType) -> N table = pa.table({"c0": pa.array([], type=arrow_type)}) with pytest.raises(ValueError): _ = vortex.array(table) + + +@pytest.mark.parametrize( + ("chunks", "arrow_type"), + [ + ([[1, 2], [3, None]], None), + ([["a", "b"], ["c"]], pa.string()), + ([[{"x": 1}, {"x": 2}], [{"x": 3}]], None), + ], +) +def test_chunked_array_combine_chunks(chunks: list[list[object]], arrow_type: pa.DataType | None) -> None: + arr = vortex.array(pa.chunked_array(chunks)) + assert isinstance(arr, vortex.ChunkedArray) + + chunked = arr.to_arrow_array(arrow_type=arrow_type) + assert isinstance(chunked, pa.ChunkedArray) + assert chunked.num_chunks == len(chunks) + + combined = arr.to_arrow_array(arrow_type=arrow_type, combine_chunks=True) + assert isinstance(combined, pa.Array) + assert combined.equals(chunked.combine_chunks()) + + +def test_chunked_struct_to_arrow_table_combine_chunks() -> None: + arr = vortex.array(pa.chunked_array([[{"x": 1}], [{"x": 2}, {"x": 3}]])) + table = arr.to_arrow_table(combine_chunks=True) + assert table.column("x").num_chunks == 1 + assert table.column("x").to_pylist() == [1, 2, 3]