apache / apache/datafusion

Incorrect Schema Adaption for CSV

Open
#4,918 6 comments 0 reactions 0 assignees View on GitHub
bug
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

**Describe the bug**

The schema adaption logic added in #1709 misbehaves for CSV data. In particular it incorrectly assumes that it can create a schema for the entire dataset that is a superset of those of the individual files, and that the CSV reader will pad any missing columns with nulls, and reorder those that appear in a different order.

In reality the CSV reader does not handle missing or reordered columns, it only accidentally works when the file schema happens to be an exact prefix of the aggregate schema. This was relying on an accidental quirk of the arrow reader prior to 30.0.0, after this the arrow reader returns an error as the schema does not match.

**To Reproduce**

Both of the following tests fail on current master

```
#[tokio::test]
async fn csv_schema_reordered() -> Result<()> {
use object_store::path::Path;

let session_ctx = SessionContext::new();

let store = InMemory::new();

let data = bytes::Bytes::from("a,b\n1,2\n3,4");
store.put(&Path::from("a.csv"), data).await.unwrap();

let data = bytes::Bytes::from("b,a\n1,2\n3,4");
store.put(&Path::from("b.csv"), data).await.unwrap();

session_ctx
.runtime_env()
.register_object_store("memory", "", Arc::new(store));

let df = session_ctx
.read_csv("memory:///", CsvReadOptions::new())
.await
.unwrap();
let result = df.collect().await.unwrap();

let expected = vec![
"+---+---+",
"| a | b |",
"+---+---+",
"| 1 | 2 |",
"| 2 | 1 |",
"| 3 | 4 |",
"| 4 | 3 |",
"+---+---+",
];

crate::assert_batches_eq!(expected, &result);

Ok(())
}
```

```
#[tokio::test]
async fn csv_schema_extra_column() -> Result<()> {
use object_store::path::Path;

let session_ctx = SessionContext::new();

let store = InMemory::new();

let data = bytes::Bytes::from("a,b\n1,2\n3,4");
store.put(&Path::from("a.csv"), data).await.unwrap();

let data = bytes::Bytes::from("a,c\n5,6\n7,8");
store.put(&Path::from("b.csv"), data).await.unwrap();

session_ctx
.runtime_env()
.register_object_store("memory", "", Arc::new(store));

let df = session_ctx
.read_csv("memory:///", CsvReadOptions::new())
.await
.unwrap();
let result = df.collect().await.unwrap();

let expected = vec![
"+---+---+---+",
"| a | b | c |",
"+---+---+---+",
"| 1 | 2 | |",
"| 3 | 4 | |",
"| 5 | | 6 |",
"| 7 | | 8 |",
"+---+---+---+",
];

crate::assert_batches_eq!(expected, &result);

Ok(())
}
```

**Expected behavior**

I think both of the following would be valid:

* Don't perform schema adaption for CSV, as they aren't a self-describing format like JSON or parquet, instead returning an error if the schema don't match
* Correctly infer the schema on a per-file basis, and use this when reading

**Additional context**
https://github.com/apache/arrow-datafusion/pull/4818#discussion_r1070588199

Contributor guide

Open the contributing guide

Research direction

Start with the csv_schema_reordered and csv_schema_extra_column reproductions using SessionContext::read_csv and inspect the CSV schema adaptation path. Determine whether CSV should reject mismatched schemas or infer each file independently, then make the reported tests pass and verify the expected batches for reordered, missing, and extra columns.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
data-engineering
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.