Skip to content

DataFusion: aggregating a dictionary column across partitions overflows the key type ("Dictionary key bigger than the key type") #220

Description

@alxmrs

Tracking issue for an upstream DataFusion bug (to file against apache/datafusion). One of two blockers for dictionary-encoding coordinate columns (#217).

Summary

GROUP BY over a dictionary-typed column fails when the combined dictionary across partitions exceeds the per-batch key width:

DataFusion error: Arrow error: Dictionary key bigger than the key type

Each input partition is a valid Dictionary(Int8, Int64) — its dictionary has ≤128 distinct values, so an Int8 key is legal per batch. But the partitions carry disjoint values, so when the aggregate combines them the combined dictionary exceeds 128 entries and the Int8 key overflows.

Minimal repro

Pure datafusion + pyarrow, no xarray-sql, no network: repros/datafusion/dict_key_overflow.py (branch claude/datafusion-upstream-repros-fs1bqv; a paste-in Rust test is in the README).

# 50 partitions, each Dictionary(Int8, Int64) of 100 DISJOINT values
#   -> combined cardinality 5000 >> Int8 max (127)
ctx.sql("SELECT k, SUM(v) FROM t GROUP BY k")  # -> "Dictionary key bigger than the key type"

The disjoint values matter: when every partition shares the same dictionary values, arrow unifies them and there is no overflow — which is why this is intermittent and version-sensitive in the wild (it does not reproduce on arrow-rs 58.3 when the values coincide).

Expected

Combining dictionary-typed columns across partitions should not fail on a key type that was valid for each input batch. It should widen the key (Int8 → Int16 → …) or decode, since a producer cannot know per batch how large the combined dictionary will become downstream.

Environment

datafusion 54.0.0 / pyarrow 23.0.0 (arrow-rs 58.3.0).

Impact on xarray-sql

Blocks dictionary-encoding coordinate columns (#217): an unchunked coordinate repeated across partitions, or a chunked coordinate whose partitions hold disjoint slices, can exceed a narrow key under streaming aggregation. Reported downstream in #217 by @ghostiee-11.

Activity

  1. alxmrs commented on Jul 3, 2026

    @alxmrs
    MemberAuthor

    Paste-ready upstream draft for apache/datafusion (title + body, AI-assistance disclosure, inline Python + Rust repro, links back here): repros/datafusion/upstream_issue_1_dict_key_overflow.md (branch claude/datafusion-upstream-repros-fs1bqv).

    Couldn't file it on apache/datafusion from this session — GitHub access here is scoped to xqlsystems/xarray-sql. Update this thread with the upstream issue link once it's filed.


    Generated by Claude Code

  2. ghostiee-11 commented on Jul 4, 2026

    @ghostiee-11
    Contributor

    Confirmed the root cause from source (arrow-rs 58.3.0 + datafusion 54.0.0, read the code; build was disk-constrained so this is by inspection, backed by the pip-repro).

    It is not an arrow-rs bug. arrow::compute::concat is type-preserving by contract (concat.rs:439 requires one shared data type), so when the merged dictionary exceeds the key type it correctly errors in merge_dictionary_values (dictionary.rs:283) rather than silently widening. The fix belongs in datafusion.

    Mechanism in datafusion 54:

    • Aggregate::group_fields (physical-plan/src/aggregates/mod.rs:414) keeps the dictionary type for a dictionary group key.
    • GroupValuesRows::emit decodes the keys via the RowConverter, then dictionary_encode_if_necessary (group_values/row.rs:296) re-encodes them back into the narrow dictionary because the group schema says so.
    • Across the RepartitionExec: Hash / FinalPartitioned boundary the combined key set exceeds Int8, so the re-encode overflows.

    #8291 (2023) decoded dict group keys to the value type and avoided exactly this; that behavior is gone in 54. apache/datafusion#7647 tracks the design and apache/datafusion#21765 (open) is moving the other way (keep dicts, optimize), without addressing the combined-cardinality overflow. Writing up the mechanism there with the minimal repro and asking whether they want decode vs widen before sending a PR. Will link it here.

  3. alxmrs commented on Jul 5, 2026

    @alxmrs
    MemberAuthor
  4. alxmrs commented on Jul 5, 2026

    @alxmrs
    MemberAuthor
  5. Rich-T-kid commented on Jul 6, 2026

    @Rich-T-kid

    @alxmrs @ghostiee-11 This issue upstream may be interesting to you apache/datafusion#23127

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

Metadata

Metadata

Assignees

No one assigned

    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