Repository navigation
Performance: Add "read strings as binary" option for parquet #12788
Description
Activity
- addedenhancementNew feature or requestNew feature or requesthelp wantedExtra attention is neededExtra attention is needed
on Oct 7, 2024 - changed the title
[-]Add "read strings as binary" option for parquet[/-][+]Performance: Add "read strings as binary" option for parquet[/+]on Oct 7, 2024 I want to try this issue, but I'm unsure if it's part of #12777. Is @jayzhan211 currently working on it?
Reacted by Austin LiuI want to try this issue, but I'm unsure if it's part of #12777. Is @jayzhan211 currently working on it?
Thank you @goldmedal
I think #12777 is version of this issue, but it changes the schema for all files, not just based on a config option
Thanks, @alamb
I see. I think I can help by doing some research on this version (maybe do a POC? 🤔 ).Reacted by Andrew LambI think a POC would be wonderful
Reacted by Jax LiuI want to try this issue, but I'm unsure if it's part of #12777. Is @jayzhan211 currently working on it?
You can take it if you would like to
Reacted by Jax Liu and Andrew Lambtake
Reacted by Andrew Lamb@alamb @jayzhan211
I drafted a PR #12816 for a simple POC. In this PR, we can use it likelet ctx = SessionContext::new(); ctx.sql( r#" CREATE EXTERNAL TABLE hits STORED AS PARQUET LOCATION 'benchmarks/data/hits_partitioned' OPTIONS ('binary_as_string' 'true') "#, ).await?.show().await?; ctx.sql("describe hits").await?.show().await?; ctx.sql(r#"select "Title" from hits limit 1"#).await?.show().await?;
The result is
+-----------------------+-----------+-------------+ | column_name | data_type | is_nullable | +-----------------------+-----------+-------------+ | WatchID | Int64 | YES | | JavaEnable | Int16 | YES | | Title | Utf8 | YES | | GoodEvent | Int16 | YES | ... +--------------------------+ | arrow_typeof(hits.Title) | +--------------------------+ | Utf8 | +--------------------------+If you want to do some experiments, this PR is easy to use.
It seems that I need to fix some CI fails, but the basic function works fine, I guess.
I'll add some tests and documents soon.Related Issue
By the way, I found an issue about casting
BinarytoStringViewwhen I tired to use thisbinary_as_stringwithschema_force_view_types.let ctx = SessionContext::new(); ctx.sql( r#" CREATE EXTERNAL TABLE hits STORED AS PARQUET LOCATION '/Users/jax/git/datafusion/benchmarks/data/hits_partitioned' OPTIONS ('binary_as_string' 'true', 'schema_force_view_types' 'true') "#, ).await?.show().await?; ctx.sql("describe hits").await?.show().await?; ctx.sql(r#"select "Title" from hits limit 1"#).await?.show().await?; ---- Error: Error during planning: Cannot cast file schema field Title of type Binary to table schema field of type Utf8View
It can be reproduced by
> select arrow_cast(arrow_cast('abc', 'Binary'), 'Utf8View'); This feature is not implemented: Unsupported CAST from Binary to Utf8View > select arrow_cast(arrow_cast('abc', 'Binary'), 'Utf8'); +-----------------------------------------------------------------+ | arrow_cast(arrow_cast(Utf8("abc"),Utf8("Binary")),Utf8("Utf8")) | +-----------------------------------------------------------------+ | abc | +-----------------------------------------------------------------+ 1 row(s) fetched. Elapsed 0.071 seconds.I guess it is an issue of
arrow-castat
https://github.com/apache/arrow-rs/blob/1be268db2237b8850161f96849353eac00cb2615/arrow-cast/src/cast/mod.rs#L209-L210Maybe we should file an issue on the arrow-rs repo. 🤔
BinaryViewworks well.> select arrow_cast(arrow_cast('abc', 'BinaryView'), 'Utf8'); +---------------------------------------------------------------------+ | arrow_cast(arrow_cast(Utf8("abc"),Utf8("BinaryView")),Utf8("Utf8")) | +---------------------------------------------------------------------+ | abc | +---------------------------------------------------------------------+ 1 row(s) fetched. Elapsed 0.034 seconds. > select arrow_cast(arrow_cast('abc', 'BinaryView'), 'Utf8View'); +-------------------------------------------------------------------------+ | arrow_cast(arrow_cast(Utf8("abc"),Utf8("BinaryView")),Utf8("Utf8View")) | +-------------------------------------------------------------------------+ | abc | +-------------------------------------------------------------------------+ 1 row(s) fetched. Elapsed 0.007 seconds.Thank you @goldmedal -- I am checking it out now.
Yes, we need to support binary -> utf8view in arrow cast
Yes, we need to support binary -> utf8view in arrow cast
Filed apache/arrow-rs#6531
Reacted by Jax LiuWhat I hope / plan to do is to combine this PR with the code in #12792 and #12809 from @Rachelint and finally get a benchmark run that shows stringview speeding up ClickBench for hits_partitioned.
Stay tuned.
Reacted by Jax LiuReacted by Daniël HeresYes, we need to support binary -> utf8view in arrow cast
Casting from binary --> utf8view via
castwill work, but won't be much/any faster than it is done todayBTW thinking more about this, I do think we need to support the cast, but in this PR we should effectively change the file schema (not just the table schema) when we setup the parquet reader (specifically with
ArrowReaderOptions::with_schema)Here is the code that does so for
Utf8-->Utf8Viewdatafusion/datafusion/core/src/datasource/physical_plan/parquet/opener.rs
Lines 126 to 129 in b821929
if let Some(merged) = coerce_file_schema_to_view_type(&table_schema, &schema) { schema = Arc::new(merged); } I think we need to also allow the switch when the table schema is
Utf8Viewand the file schema is Binary/BinaryViewReacted by Jax LiuBTW thinking more about this, I do think we need to support the cast, but in this PR we should effectively change the file schema (not just the table schema) when we setup the parquet reader (specifically with
ArrowReaderOptions::with_schema)This idea sounds great. It seems if we can apply the new schema when reading file, we can save one time casting. Just read as string.
I tried to follow the implementation of StringView to apply the new schema using
with_schemabut I got casting error.Parquet error: Arrow: incompatible arrow schema, the following fields could not be cast:I can reprodcue this error on the arrow-rs side by added a test case in
parquet/src/arrow/arrow_reader/mod.rs#[test] fn test_cast_binary_utf8() { let original_fields = Fields::from(vec![ Field::new("binary_to_utf8", ArrowDataType::Binary, false), ]); let file = write_parquet_from_iter(vec![ ( "binary_to_utf8", Arc::new(BinaryArray::from(vec![b"one".as_ref(), b"two".as_ref()])) as ArrayRef, ), ]); let supplied_fields = Fields::from(vec![ Field::new("binary_to_utf8", ArrowDataType::Utf8, false), ]); let options = ArrowReaderOptions::new().with_schema(Arc::new(Schema::new(supplied_fields))); let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options( file.try_clone().unwrap(), options, ) .expect("reader builder with schema") .build() .expect("reader with schema"); let batch = arrow_reader.next().unwrap().unwrap(); assert_eq!(batch.num_columns(), 1); assert_eq!(batch.num_rows(), 2); assert_eq!( batch .column(0) .as_any() .downcast_ref::<StringArray>() .expect("downcast to string") .iter() .collect::<Vec<_>>(), vec![Some("one"), Some("two")] ); }
The output is
reader builder with schema: ArrowError("incompatible arrow schema, the following fields could not be cast: [binary_to_utf8]")I tired to fix it through adding more pattern match at
https://github.com/apache/arrow-rs/blob/5508978a3c5c4eb65ef6410e097887a8adaba38a/parquet/src/arrow/schema/primitive.rs#L40(DataType::Binary, DataType::Utf8) => hint,
It can work well but I'm not pretty sure if this way makes sense 🤔
BTW, this way might be not safe. If the data isn't a valid Utf8 binary like
Arc::new(BinaryArray::from(vec![b"\xDE\x00\xFF".as_ref()])) as ArrayRef,This casting will fail by
InvalidArgumentError("Invalid UTF8 sequence at string index 0 (0..3): invalid utf-8 sequence of 1 bytes from index 0")Maybe we also need to make this behaivor be optional on the arrow-rs side? 🤔
I tried to follow the implementation of StringView to apply the new schema using
with_schemabut I got casting error.Parquet error: Arrow: incompatible arrow schema, the following fields could not be cast:Submitted apache/arrow-rs#6539 for it.
Submitted apache/arrow-rs#6539 for it.
Thank you @goldmedal -- that looks great.
TLDR I would like to add a new
binary_as_stringoption for paruetIs your feature request related to a problem or challenge?
The Real Problem
The primary problem is that the ClickBench queries slow down when we enable StringView by default only for the
hits_partitionedversion of the datasetOne of the reasons is that reading a column as a
BinaryViewArrayand then casting toUtf8ViewArrayis significantly slower than reading the data from parquet as aUtf8ViewArray(due to optimizations in the parquet decoder).This is not a problem with
StringArray-->BinaryArraybecause reading a column as aBinaryArrayand then casting toUtf8Arrayis about the same speed as reading as Utf8ArrayThe core issue is that for
hits_partitionedthe "String" columns in the schema are marked as binary (not Utf8) and thus the slower conversion path is used.The background:
hits_partitionedhas "string" columns marked as "Binary"Clickbench has 2 versions of the parquet dataset (see docs here)
hits.parquet(a single 14G parquet file)athena_partitioned/hits_{0..99}.parquetHowever, the SCHEMA is different between these two files
hits.parquethasStrings:DataFusion recognizes this as Utf8
hits_partitionedhas the string columns asBinary:Which datafusion correctly interprets as
Binary:Describe the solution you'd like
I would like a way to treat binary columns in the hits_partitioned dataset as Strings.
This is the right thing to do for the
hits_partitioneddataset, but I am not sure it is the right thing to do for all files, so I think we need some flagDescribe alternatives you've considered
I propose adding a
binary_as_stringoption to the parquet reader likeIn fact, it seems as if duckDB does exactly this note the (
binary_as_string=True) item herehttps://duckdb.org/docs/data/parquet/overview.html#parameters
And it is set in clickbench scripts:
https://github.com/ClickHouse/ClickBench/blob/a6615de8e6eae45f8c1b047df0073fe32f43499f/duckdb-parquet/create.sql#L6
Additional context
Trino also explictly specifies the types of files in the
hits_partitioneddatasethttps://github.com/ClickHouse/ClickBench/blob/a6615de8e6eae45f8c1b047df0073fe32f43499f/trino/create_partitioned.sql#L1-L107
We could argue that adding something just for benchmarking goes against the spirit of the test, however, I think we can make a reasonable argument on correctness grounds too
For example, if you run this query today
You get this arguably nonsensical result (the search phrase is shown as binary):
However if you run the same query with the proper schema (strings) you see a real search phrase