apache / apache/datafusion

Consolidate datafusion, arrow-cpp, and substrait's handling of non-substrait arrow types

Open
#12,181 2 comments 1 reaction 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

I am trying to take a filter expression created by pyarrow and convert it into a filter expression for Datafusion to satisfy. I am using Substrait to do this. Everything works fine when I use the standard Substrait types. However, when I use normal Arrow types that are not Substrait types (e.g. unsigned integers, large containers) I run into problems.

It seems that arrow-cpp (admittedly, me, in this case) and datafusion have taken different approaches to handling these limitations.

In arrow-cpp the types that expand or change the valid range of values (e.g. unsigned integers, large containers) are converted to extension types. This process is documented in https://github.com/apache/arrow/blob/main/format/substrait/extension_types.yaml

In datafusion it appears these types are expected to use the nearest substrait match (e.g. signed integer, small container) with a type variation.

### Describe the solution you'd like

I am admittedly biased (given I implemented one of the two disagreeing components) but I favor the extension types approach. Type variations are defined in Substrait as this:

> Type variations may be used to represent differences in representation between different consumers. For example, an engine might support dictionary encoding for a string, or could be using either a row-wise or columnar representation of a struct. All variations of a type are expected to have the same semantics when operated on by functions or other expressions.

Given that definition, I do not think it is valid to say that an unsigned integer is a variation of a signed integer (they do not have the same outputs for all functions). I do believe things like the view types and dictionary encoding are valid type variations.

### Describe alternatives you've considered

The alternative would be to change arrow-cpp to also use type variations. Though I'd like some consensus from the Substrait community that this is a valid use of type variations before taking that approach.

At the moment I am working around this issue by simply removing any non-standard types from the input schema (this works as long as the filter isn't referencing those types).

### Additional context

_No response_

Contributor guide

Open the contributing guide

Research direction

Start by reading the Substrait extension_types.yaml guidance and comparing how arrow-cpp and DataFusion represent unsigned integers and large containers. Review the existing discussion for agreement on extension types versus type variations; done means the projects have a decided, consistent handling for these non-Substrait Arrow types.

Written by the indexing model from the issue text.

Assessment

Tech stack
python, rust
Domain
data-engineering
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.