apache / apache/arrow-rs

Add coerce_types flag to parquet ArrowWriter

Open
#1,938 12 comments 0 reactions 1 assignee Claimed by @CuteChuanChuan View on GitHub
enhancement help wanted
Dominant language
Rust
Stars
3.6k
Forks
1.3k
Avg merge
2d 18h
Merged PRs (30d)
169

Description

**Is your feature request related to a problem or challenge? Please describe what you are trying to do.**

As discussed in #1666 not all types can be represented within a parquet schema.

**Describe the solution you'd like**

The consensus appears to be to:

* By default faithfully round-trip the source data, performing no potentially lossy type conversion
* Add a coerce_types flag that will use the arrow cast kernels to coerce incompatible types prior to writing them

In particular

### Date64

If not coerce_types, write as Int64 and embed logical type in arrow schema only. Otherwise case to Date32

### Timestamp

If not coerce_types, write as is, setting LogicalType / ConvertedType only where appropriate.

If coerce_types, cast to a UTC timestamp with the closest supported time unit, likely needing #1936.

### Interval

If not coerce_types, write as FixedSizeBinaryArray matching the arrow representation and store logical type in arrow schema.

If coerce_types, convert to the relevant parquet representation.

**Describe alternatives you've considered**

See #1666

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.