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
112 changes: 108 additions & 4 deletions python/pyarrow/src/arrow/python/arrow_to_pandas.cc
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
#include <cmath>
#include <cstdint>
#include <iostream>
#include <limits>
#include <memory>
#include <mutex>
#include <string>
Expand Down Expand Up @@ -747,6 +748,64 @@ Status ConvertStruct(PandasOptions options, const ChunkedArray& data,
return Status::OK();
}

// Whether decoding `arr` into its string/binary value type would overflow the
// 32-bit offsets of the dense array (GH-50842), i.e. whether Take would reject
// it. The length of the dense data is computed from the dictionary offsets.
bool DecodingWouldOverflow(const DictionaryArray& arr) {
// As in the Take kernel for binary-like arrays
constexpr int64_t kOffsetLimit = std::numeric_limits<int32_t>::max() - 1;
const auto& dictionary = checked_cast<const BinaryArray&>(*arr.dictionary());
const int64_t dictionary_length = dictionary.length();
// Cheap bounds first: no overflow is possible if every row referred to the
// whole dictionary, or if no value is longer than kOffsetLimit / length
const int64_t length = arr.length();
const int64_t dictionary_bytes =
dictionary.value_offset(dictionary_length) - dictionary.value_offset(0);
if (dictionary_bytes <= 0 || length <= kOffsetLimit / dictionary_bytes) {
return false;
}
if (dictionary_length <= length) {
const int64_t max_value_length = kOffsetLimit / length;
int64_t i = 0;
while (i < dictionary_length && dictionary.value_length(i) <= max_value_length) {
++i;
}
if (i == dictionary_length) {
return false;
}
}
const bool has_nulls = arr.null_count() > 0;
int64_t total_length = 0;
for (int64_t i = 0; i < length; ++i) {
if (has_nulls && arr.IsNull(i)) {
continue;
}
const int64_t index = arr.GetValueIndex(i);
if (index < 0 || index >= dictionary_length) {
// Let the decoding report the invalid index
return false;
}
total_length += dictionary.value_length(index);
if (total_length > kOffsetLimit) {
return true;
}
}
return false;
}

bool DecodingWouldOverflow(const ChunkedArray& arr) {
const auto& value_type = checked_cast<const DictionaryType&>(*arr.type()).value_type();
if (!is_binary_like(value_type->id())) {
return false;
}
for (int c = 0; c < arr.num_chunks(); c++) {
if (DecodingWouldOverflow(checked_cast<const DictionaryArray&>(*arr.chunk(c)))) {
return true;
}
}
return false;
}

Status DecodeDictionaries(MemoryPool* pool, const std::shared_ptr<DataType>& dense_type,
ArrayVector* arrays) {
compute::ExecContext ctx(pool);
Expand Down Expand Up @@ -1380,6 +1439,40 @@ struct ObjectWriterVisitor {
return ConvertStruct(options, data, out_values);
}

// Dictionaries whose decoding would overflow the 32-bit offsets of the dense
// array (GH-50842) are not decoded. Instead, the dictionary values are
// converted once and referenced for each index.
Status Visit(const DictionaryType& type) {
for (int c = 0; c < data.num_chunks(); c++) {
const auto& arr = checked_cast<const DictionaryArray&>(*data.chunk(c));
const int64_t dictionary_length = arr.dictionary()->length();
std::vector<PyObject*> values(dictionary_length, nullptr);
ChunkedArray dictionary(arr.dictionary());
ObjectWriterVisitor values_visitor{options, dictionary, values.data()};
Status st = VisitTypeInline(*arr.dictionary()->type(), &values_visitor);
for (int64_t i = 0; st.ok() && i < arr.length(); ++i) {
if (arr.IsNull(i)) {
Py_INCREF(Py_None);
*out_values = Py_None;
} else {
const int64_t index = arr.GetValueIndex(i);
if (index < 0 || index >= dictionary_length) {
st = Status::IndexError("Index ", index, " out of bounds");
break;
}
Py_INCREF(values[index]);
*out_values = values[index];
}
++out_values;
}
for (PyObject* value : values) {
Py_XDECREF(value);
}
RETURN_NOT_OK(st);
}
return Status::OK();
}

Status Visit(const ExtensionType& type) {
if (!IsUuidExtension(type)) {
return Status::NotImplemented("No implemented conversion to object dtype: ",
Expand Down Expand Up @@ -2275,9 +2368,16 @@ static Status GetPandasWriterType(const ChunkedArray& data, const PandasOptions&
}
*output_type = PandasWriter::OBJECT;
} break;
case Type::DICTIONARY:
*output_type = PandasWriter::CATEGORICAL;
break;
case Type::DICTIONARY: {
// A string/binary dictionary that is not decoded because the dense array
// would overflow is converted to objects instead, see
// ObjectWriterVisitor::Visit(const DictionaryType&)
const auto& value_type =
checked_cast<const DictionaryType&>(*data.type()).value_type();
*output_type = options.decode_dictionaries && is_binary_like(value_type->id())
? PandasWriter::OBJECT
: PandasWriter::CATEGORICAL;
} break;
case Type::EXTENSION:
// UUID has a native object conversion to uuid.UUID. Other extension
// types continue through the pandas ExtensionArray protocol.
Expand Down Expand Up @@ -2595,7 +2695,11 @@ Status ConvertArrayToPandas(const PandasOptions& options, std::shared_ptr<Array>
Status ConvertChunkedArrayToPandas(const PandasOptions& options,
std::shared_ptr<ChunkedArray> arr, PyObject* py_ref,
PyObject** out) {
if (options.decode_dictionaries && arr->type()->id() == Type::DICTIONARY) {
// Decoding a string/binary dictionary can overflow the 32-bit offsets of
// the dense array (GH-50842): such arrays are converted to objects without
// decoding instead, see ObjectWriterVisitor::Visit(const DictionaryType&)
if (options.decode_dictionaries && arr->type()->id() == Type::DICTIONARY &&
!DecodingWouldOverflow(*arr)) {
// XXX we should return an error as below if options.zero_copy_only
// is true, but that would break compatibility with existing tests.
const auto& dense_type =
Expand Down
32 changes: 32 additions & 0 deletions python/pyarrow/tests/test_array.py
Original file line number Diff line number Diff line change
Expand Up @@ -1082,6 +1082,38 @@ def test_dictionary_to_numpy():
)


@pytest.mark.numpy
def test_dictionary_to_numpy_without_decoding():
# GH-50842: a dictionary whose decoded form would overflow the 32-bit
# offsets of its value type (16385 * 2**17 > INT32_MAX bytes here) is
# converted to objects without decoding
value = "a" * 2**17
arr = pa.DictionaryArray.from_arrays(
np.zeros(16385, dtype=np.int16), pa.array([value]))
result = np.asarray(arr)
assert len(result) == len(arr)
assert result[0] == value
assert result[-1] == value

# With a null index (row 2) and a null dictionary value (row 1)
indices = np.zeros(16387, dtype=np.int16)
indices[1] = 1
arr = pa.DictionaryArray.from_arrays(
pa.array(indices, mask=np.arange(16387) == 2), pa.array([value, None]))
result = arr.to_numpy(zero_copy_only=False)
assert result[0] == value
assert result[1] is None
assert result[2] is None
assert result[-1] == value

# Binary values
arr = pa.DictionaryArray.from_arrays(
np.zeros(16385, dtype=np.int16), pa.array([value.encode()]))
result = np.asarray(arr)
assert result[0] == value.encode()
assert result[-1] == value.encode()


@pytest.mark.numpy
def test_dictionary_from_boxed_arrays():
indices = np.repeat([0, 1, 2], 2)
Expand Down
18 changes: 18 additions & 0 deletions python/pyarrow/tests/test_pandas.py
Original file line number Diff line number Diff line change
Expand Up @@ -2518,6 +2518,24 @@ def test_list_of_dictionary(self):
expected[2] = None
tm.assert_series_equal(arr.to_pandas(), expected)

def test_list_of_dictionary_without_decoding(self):
# GH-50842: the dictionary child is converted to objects without
# being decoded into a dense array, which would overflow its 32-bit
# offsets (16385 * 2**17 > INT32_MAX bytes here)
value = "a" * 2**17
child = pa.DictionaryArray.from_arrays(
np.zeros(16385, dtype=np.int16), pa.array([value]))
arr = pa.ListArray.from_arrays([0, len(child)], child)
result = arr.to_pandas()
assert len(result[0]) == len(child)
assert result[0][0] == value

# Also with several chunks
result = pa.chunked_array([arr, arr]).to_pandas()
assert len(result) == 2
assert len(result[1]) == len(child)
assert result[1][-1] == value

@pytest.mark.large_memory
def test_auto_chunking_on_list_overflow(self):
# ARROW-9976
Expand Down
Loading