Skip to content

[Python] Support DictionaryArray -> numpy/pandas conversion when the decoded data exceeds the 32-bit offset limit #50842

Description

@pearu

Describe the enhancement requested

Once #50840 is fixed (#50841), converting a dictionary array whose decoded form exceeds INT32_MAX bytes raises cleanly instead of crashing:

import numpy as np, pyarrow as pa

arr = pa.DictionaryArray.from_arrays(
    np.zeros(50_000_000, dtype=np.int16),
    pa.array(["a" * 50]),
)
np.asarray(arr)
# pyarrow.lib.ArrowInvalid: Take operation overflowed binary array capacity

That is a strict improvement over a segfault, but the conversion still fails on data that is entirely representable in the output. The result of np.asarray here is a numpy object array of Python str — a format with no 2 GiB limit. The failure comes purely from an intermediate step.

Why it fails

Array.to_numpy(zero_copy_only=False) sets decode_dictionaries=True. In arrow_to_pandas.cc, ConvertChunkedArrayToPandas() then decodes by casting to the dictionary's value type:

const auto& dense_type =
    checked_cast<const DictionaryType&>(*arr->type()).value_type();   // string, int32 offsets
RETURN_NOT_OK(DecodeDictionaries(options.pool, dense_type, &arr));

So a 50M × 50-byte decode has to fit in a 32-bit-offset string array. It cannot — hence the error. But that dense array is a pure implementation artifact; nothing downstream needs it to be a string array specifically.

Approach (a): widen the intermediate type

Use large_string/large_binary as the decode target instead of string/binary.

Pros

  • Very small, local change — essentially picking a different dense_type
  • Low risk; reuses the existing cast/take machinery unchanged
  • Invisible to users: the final output is an object array either way
  • Fixes the to_numpy and to_pandas paths together

Cons

  • Still materializes the entire dense array (2.5 GB in the example) purely as scratch
  • Still creates one Python object per element — 50M separate str objects for 50M indices — so peak memory stays very large even though the input holds a single distinct value
  • Slightly increases memory for the common, non-overflowing case (8-byte offsets), unless a size pre-check is added to decide when to promote
  • Treats the symptom rather than the redundancy

Approach (b): skip the dense intermediate entirely

For the object-output path, convert each dictionary value to a PyObject once, then walk the indices and Py_INCREF the corresponding object into the output array.

Pros

  • No dense intermediate at all — the 2.5 GB scratch buffer disappears
  • For the example, one Python string is allocated instead of 50M; memory and time drop by orders of magnitude
  • Removes the size ceiling entirely rather than raising it
  • Also speeds up the ordinary, non-overflowing dictionary → object conversion, which is the common categorical/pandas path
  • Output is semantically identical: object arrays hold references, and Python strings are immutable, so sharing is unobservable

Cons

  • Larger change: needs a dedicated dictionary branch in the object-writer instead of reusing DecodeDictionaries
  • Null indices and index bounds must be handled explicitly rather than inherited from the cast
  • Only covers the object-output path; any path that genuinely needs a dense Arrow array (e.g. strings_to_categorical) would still want (a)
  • Callers relying on distinct object identity per element would see shared references — not a documented guarantee, and already untrue for interned strings, but worth noting

Recommendation

(b), because it addresses the actual redundancy rather than raising a ceiling. Dictionary-encoded data is used precisely when values repeat, so materializing one Python object per index rather than per dictionary entry is wasted work in every case — the overflow is just where it becomes fatal. It also improves the common path, not only the pathological one.

(a) remains worth keeping in reserve for any remaining code path that must produce a dense array; the two are complementary rather than exclusive.

I'm happy to implement (b), pending agreement on the direction.

Component(s)

Python


🤖 Drafted by Claude Code (an AI agent) and reviewed & approved by pearu.

Activity

  1. KHARSHAVARDHAN-eng commented on Aug 13, 2026

    @KHARSHAVARDHAN-eng

    @pearu ...I’d be interested in contributing to this. The direct dictionary → Python object approach (b) makes sense to me, especially since it avoids the large intermediate allocation. If you’re not already planning to implement it, I’d be happy to work on the implementation and regression tests.

  2. pearu commented on Aug 13, 2026

    @pearu
    ContributorAuthor

    Thanks for the offer, and good to have a second voice for (b) — that's the direction I'd settled on too.

    I am planning to implement it, so I'll take this one. The prerequisite fix (#50841) has landed, so the path is clear.

    If you'd like something adjacent to pick up: (a) is still potentially worth doing as a complement, for code paths that genuinely need a dense array and so aren't covered by (b) — strings_to_categorical being the obvious one. That's separable from this issue. Otherwise, review on the PR once it's up would be very welcome.

    Worth noting no maintainer has commented on the direction yet, so (b) is still two contributors' preference rather than a settled decision — input from an Arrow maintainer would be welcome before I get too far in.


    🤖 Drafted by Claude Code (an AI agent) and reviewed & approved by pearu.

  3. pearu commented on Oct 9, 2026

    @pearu
    ContributorAuthor

    @rok @jorisvandenbossche following the suggestion at the dev sync, I looked at whether the deduplicate_objects machinery could solve this. It can't as it stands, for two reasons:

    1. It runs too late. ConvertChunkedArrayToPandas first decodes the dictionary into a dense Arrow array (DecodeDictionaries → compute::Cast(dictionary → value_type), which is a take of the dictionary by the indices). Only that dense string array then reaches the object writer, ConvertAsPyObjects, where deduplicate_objects lives. The decode is exactly the step that overflows the 32-bit offsets, so deduplication never gets to run.
    2. On the np.asarray / to_numpy path it isn't enabled at all. to_numpy sets decode_dictionaries, zero_copy_only and to_numpy, and leaves deduplicate_objects at its C++ default of false. (to_pandas enables it, but reason 1 still applies.)

    Approach: let the object writer consume the dictionary array directly instead of a decoded copy: convert each dictionary value to a Python object once, then fill the output by Py_INCREF-ing the object for each index. This is the unique_values + Py_INCREF part of deduplicate_objects, minus the memo table, because the dictionary already is the set of unique values.

    • Reason 1 goes away because there is no dense intermediate any more; nothing before the object stage can overflow.
    • Reason 2 goes away because the dictionary path doesn't depend on deduplicate_objects being enabled; it applies whenever a binary-like dictionary is decoded into an object array, which covers to_numpy as is.

    Plan: one PR in arrow_to_pandas.cc, with regression tests for the two entry points that hit the decode: to_numpy/np.asarray, and to_pandas of nested dictionaries (e.g. list<dictionary<string>>).

    I'll prepare the PR unless there are objections to the direction.

  4. jorisvandenbossche commented on Oct 9, 2026

    @jorisvandenbossche
    Member

    Yes, that sounds good!

  5. added a commit that references this issue on Oct 9, 2026
  6. pearu commented on Oct 9, 2026

    @pearu
    ContributorAuthor

    Follow-up on the plan above: the PR (#52547) narrows it. Benchmarking the always-on variant (dictionary path for every string/binary dictionary) against the current decode showed that it is faster and far lighter for ordinary dictionaries (10²–10⁴ entries: 2–4× faster, ~100× less memory), but worse in three situations: dictionaries beyond ~10⁵ entries (the table of Python objects no longer fits in cache, up to 3× slower), many chunks sharing one large dictionary (the dictionary is converted once per chunk; 1000 chunks × 10³ rows over a 10⁵-entry dictionary took 10 s instead of ~50 ms), and small slices of arrays with large dictionaries.

    So the PR takes the dictionary path only when decoding would actually overflow: a cheap check on the dictionary offsets and indices predicts whether Take would reject the decode, and everything else goes through the existing code unchanged (measured at parity on 10⁶–10⁷ rows with 10²–10⁶-entry dictionaries). Using the dictionary path more widely, where it is the better choice, can be a separate change.


    🤖 Drafted by Claude Code (an AI agent) and reviewed & approved by pearu.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

No labels
No labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions