Skip to main content

parquet/
variant.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! ⚠️ Experimental Support for reading and writing [`Variant`]s to / from Parquet files ⚠️
19//!
20//! This is a 🚧 Work In Progress
21//!
22//! Note: Requires the `variant_experimental` feature of the `parquet` crate to be enabled.
23//!
24//! # Features
25//! * Representation of [`Variant`], and [`VariantArray`] for working with
26//!   Variant values (see [`parquet_variant`] for more details)
27//! * Kernels for working with arrays of Variant values
28//!   such as conversion between `Variant` and JSON, and shredding/unshredding
29//!   (see [`parquet_variant_compute`] for more details)
30//!
31//! # Example: Writing a Parquet file with Variant column
32//! ```rust
33//! # use parquet::variant::{VariantArray, VariantType, VariantArrayBuilder, VariantBuilderExt};
34//! # use std::sync::Arc;
35//! # use arrow_array::{Array, ArrayRef, RecordBatch};
36//! # use arrow_schema::{DataType, Field, Schema};
37//! # use parquet::arrow::ArrowWriter;
38//! # fn main() -> Result<(), parquet::errors::ParquetError> {
39//!  // Use the VariantArrayBuilder to build a VariantArray
40//!  let mut builder = VariantArrayBuilder::new(3);
41//!  builder.new_object().with_field("name", "Alice").finish(); // row 1: {"name": "Alice"}
42//!  builder.append_value("such wow"); // row 2: "such wow" (a string)
43//!  let array = builder.build();
44//!
45//!  // Since VariantArray is an ExtensionType, it needs to be converted
46//!  // to an ArrayRef and Field with the appropriate metadata
47//!  // before it can be written to a Parquet file
48//!  let field = array.field("data");
49//!  let array = ArrayRef::from(array);
50//!  // create a RecordBatch with the VariantArray
51//!  let schema = Schema::new(vec![field]);
52//!  let batch = RecordBatch::try_new(Arc::new(schema), vec![array])?;
53//!
54//!  // Now you can write the RecordBatch to the Parquet file, as normal
55//!  let file = std::fs::File::create("variant.parquet")?;
56//!  let mut writer = ArrowWriter::try_new(file, batch.schema(), None)?;
57//!  writer.write(&batch)?;
58//!  writer.close()?;
59//!
60//! # std::fs::remove_file("variant.parquet")?;
61//! # Ok(())
62//! # }
63//! ```
64//!
65//! # Example: Writing JSON into a Parquet file with Variant column
66//! ```rust
67//! # use std::sync::Arc;
68//! # use arrow_array::{ArrayRef, RecordBatch, StringArray};
69//! # use arrow_schema::Schema;
70//! # use parquet::variant::{json_to_variant, VariantArray};
71//! # use parquet::arrow::ArrowWriter;
72//! # fn main() -> Result<(), parquet::errors::ParquetError> {
73//!  // Create an array of JSON strings, simulating a column of JSON data
74//!  let input_array: ArrayRef = Arc::new(StringArray::from(vec![
75//!   Some(r#"{"name": "Alice", "age": 30}"#),
76//!   Some(r#"{"name": "Bob", "age": 25, "address": {"city": "New York"}}"#),
77//!   None,
78//!   Some("{}"),
79//!  ]));
80//!
81//!  // Convert the JSON strings to a VariantArray
82//!  let array: VariantArray = json_to_variant(&input_array)?;
83//!  // create a RecordBatch with the VariantArray
84//!  let schema = Schema::new(vec![array.field("data")]);
85//!  let batch = RecordBatch::try_new(Arc::new(schema), vec![ArrayRef::from(array)])?;
86//!
87//!  // write the RecordBatch to a Parquet file as normal
88//!  let file = std::fs::File::create("variant-json.parquet")?;
89//!  let mut writer = ArrowWriter::try_new(file, batch.schema(), None)?;
90//!  writer.write(&batch)?;
91//!  writer.close()?;
92//! # std::fs::remove_file("variant-json.parquet")?;
93//! # Ok(())
94//! # }
95//! ```
96//!
97//! # Example: Reading a Parquet file with Variant column
98//!
99//! Use the [`VariantType`] extension type to find the Variant column:
100//!
101//! ```
102//! # use std::sync::Arc;
103//! # use std::path::PathBuf;
104//! # use arrow_array::{ArrayRef, RecordBatch, RecordBatchReader};
105//! # use parquet::variant::{Variant, VariantArray, VariantType};
106//! # use parquet::arrow::arrow_reader::ArrowReaderBuilder;
107//! # fn main() -> Result<(), parquet::errors::ParquetError> {
108//! # use arrow_array::StructArray;
109//! # fn file_path() -> PathBuf { // return a testing file path
110//! #    PathBuf::from(arrow::util::test_util::parquet_test_data())
111//! #   .join("..")
112//! #   .join("shredded_variant")
113//! #   .join("case-075.parquet")
114//! # }
115//! // Read the Parquet file using standard Arrow Parquet reader.
116//! // Note this file has 2 columns: "id", "var", and the "var" column
117//  // contains a variant that looks like this:
118//  // "Variant(metadata=VariantMetadata(dict={}), value=Variant(type=STRING, value=iceberg))"
119//! let file = std::fs::File::open(file_path())?;
120//! let mut reader = ArrowReaderBuilder::try_new(file)?.build()?;
121//!
122//! // You can check if a column contains a Variant using
123//! // the VariantType extension type
124//! let schema = reader.schema();
125//! let field = schema.field_with_name("var")?;
126//! assert!(field.has_valid_extension_type::<VariantType>());
127//!
128//! // The reader will yield RecordBatches with a StructArray
129//! // to convert them to VariantArray, use VariantArray::try_new
130//! let batch = reader.next().unwrap().unwrap();
131//!
132//! let col = batch.column_by_name("var").unwrap();
133//! let var_array = VariantArray::try_new(col)?;
134//! assert_eq!(var_array.len(), 1);
135//! let var_value: Variant = var_array.value(0);
136//! assert_eq!(var_value, Variant::from("iceberg")); // the value in case-075.parquet
137//! # Ok(())
138//! # }
139//! ```
140pub use parquet_variant::*;
141pub use parquet_variant_compute::*;
142
143#[cfg(test)]
144mod tests {
145    use crate::arrow::ArrowWriter;
146    use crate::arrow::arrow_reader::ArrowReaderBuilder;
147    use crate::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
148    use crate::file::reader::ChunkReader;
149    use arrow::util::test_util::parquet_test_data;
150    use arrow_array::{
151        Array, ArrayRef, BinaryViewArray, FixedSizeListArray, Int64Array, RecordBatch, StructArray,
152        new_null_array,
153    };
154    use arrow_schema::{DataType, Field, Fields, Schema};
155    use bytes::Bytes;
156    use parquet_variant::{EMPTY_VARIANT_METADATA_BYTES, Variant, VariantBuilderExt};
157    use parquet_variant_compute::{
158        VariantArray, VariantArrayBuilder, VariantType, unshred_variant,
159    };
160    use std::path::PathBuf;
161    use std::sync::Arc;
162
163    #[test]
164    fn roundtrip_basic() {
165        roundtrip(variant_array());
166    }
167
168    /// Ensure a file with Variant LogicalType, written by another writer in
169    /// parquet-testing, can be read as a VariantArray
170    #[test]
171    fn read_logical_type() {
172        // Note: case-075 2 columns ("id", "var")
173        // The variant looks like this:
174        // "Variant(metadata=VariantMetadata(dict={}), value=Variant(type=STRING, value=iceberg))"
175        let batch = read_shredded_variant_test_case("case-075.parquet");
176
177        assert_variant_metadata(&batch, "var");
178        let var_column = batch.column_by_name("var").expect("expected var column");
179        let var_array =
180            VariantArray::try_new(&var_column).expect("expected var column to be a VariantArray");
181
182        // verify the value
183        assert_eq!(var_array.len(), 1);
184        assert!(var_array.is_valid(0));
185        let var_value = var_array.value(0);
186        assert_eq!(var_value, Variant::from("iceberg"));
187    }
188
189    #[test]
190    fn read_fixed_size_list_typed_value_as_list() {
191        let element_values: ArrayRef = Arc::new(Int64Array::from(vec![1, 2, 3, 4]));
192        let element_value = new_null_array(&DataType::BinaryView, 4);
193        let element_fields = Fields::from(vec![
194            Field::new("value", DataType::BinaryView, true),
195            Field::new("typed_value", DataType::Int64, true),
196        ]);
197        let elements: ArrayRef = Arc::new(StructArray::new(
198            element_fields,
199            vec![element_value, element_values],
200            None,
201        ));
202        let item_field = Arc::new(Field::new("item", elements.data_type().clone(), true));
203        let typed_value: ArrayRef =
204            Arc::new(FixedSizeListArray::new(item_field, 2, elements, None));
205        let metadata: ArrayRef = Arc::new(BinaryViewArray::from_iter_values(std::iter::repeat_n(
206            EMPTY_VARIANT_METADATA_BYTES,
207            2,
208        )));
209        let value = new_null_array(&DataType::BinaryView, 2);
210        let fields = Fields::from(vec![
211            Field::new("metadata", DataType::BinaryView, false),
212            Field::new("value", DataType::BinaryView, true),
213            Field::new("typed_value", typed_value.data_type().clone(), true),
214        ]);
215        let source = StructArray::new(fields, vec![metadata, value, typed_value], None);
216        let field = Field::new("data", source.data_type().clone(), false);
217        let batch =
218            RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![Arc::new(source)])
219                .unwrap();
220
221        let buffer = write_to_buffer(&batch);
222        let result = read_to_batch(Bytes::from(buffer));
223        let column = result.column_by_name("data").unwrap();
224        let variant = VariantArray::try_new(column).unwrap();
225        assert!(matches!(
226            variant.typed_value_column().unwrap().data_type(),
227            DataType::List(_)
228        ));
229
230        let unshredded = unshred_variant(&variant).unwrap();
231        assert!(unshredded.typed_value_column().is_none());
232        assert_eq!(unshredded.len(), 2);
233    }
234
235    /// Writes a variant to a parquet file and ensures the parquet logical type
236    /// annotation is correct
237    #[test]
238    fn write_logical_type() {
239        let array = variant_array();
240        let batch = variant_array_to_batch(array);
241        let buffer = write_to_buffer(&batch);
242
243        // read the parquet file's metadata and verify the logical type
244        let metadata = read_metadata(&Bytes::from(buffer));
245        let schema = metadata.file_metadata().schema_descr();
246        let fields = schema.root_schema().get_fields();
247        assert_eq!(fields.len(), 1);
248        let field = &fields[0];
249        assert_eq!(field.name(), "data");
250        // data should have been written with the Variant logical type
251        assert_eq!(
252            field.get_basic_info().logical_type_ref(),
253            Some(&crate::basic::LogicalType::variant(None))
254        );
255    }
256
257    /// Return a VariantArray with 3 rows:
258    ///
259    /// 1. `{"name": "Alice"}`
260    /// 2. `"such wow"` (a string)
261    /// 3. `null`
262    fn variant_array() -> VariantArray {
263        let mut builder = VariantArrayBuilder::new(3);
264        // row 1: {"name": "Alice"}
265        builder.new_object().with_field("name", "Alice").finish();
266        // row 2: "such wow" (a string)
267        builder.append_value("such wow");
268        // row 3: null
269        builder.append_null();
270        builder.build()
271    }
272
273    /// Writes a VariantArray to a parquet file and reads it back, verifying that
274    /// the data is the same
275    fn roundtrip(array: VariantArray) {
276        let source_batch = variant_array_to_batch(array);
277        assert_variant_metadata(&source_batch, "data");
278
279        let buffer = write_to_buffer(&source_batch);
280        let result_batch = read_to_batch(Bytes::from(buffer));
281        assert_variant_metadata(&result_batch, "data");
282        assert_eq!(result_batch, source_batch); // NB this also checks the schemas
283    }
284
285    /// creates a RecordBatch with a single column "data" from a VariantArray,
286    fn variant_array_to_batch(array: VariantArray) -> RecordBatch {
287        let field = array.field("data");
288        let schema = Schema::new(vec![field]);
289        RecordBatch::try_new(Arc::new(schema), vec![ArrayRef::from(array)]).unwrap()
290    }
291
292    /// writes a RecordBatch to memory buffer and returns the buffer
293    fn write_to_buffer(batch: &RecordBatch) -> Vec<u8> {
294        let mut buffer = vec![];
295        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
296        writer.write(batch).unwrap();
297        writer.close().unwrap();
298        buffer
299    }
300
301    /// Reads the Parquet metadata
302    fn read_metadata<T: ChunkReader + 'static>(input: &T) -> ParquetMetaData {
303        let mut reader = ParquetMetaDataReader::new();
304        reader.try_parse(input).unwrap();
305        reader.finish().unwrap()
306    }
307
308    /// Reads a RecordBatch from a reader (e.g. Vec or File)
309    fn read_to_batch<T: ChunkReader + 'static>(reader: T) -> RecordBatch {
310        let reader = ArrowReaderBuilder::try_new(reader)
311            .unwrap()
312            .build()
313            .unwrap();
314        let mut batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().unwrap();
315        assert_eq!(batches.len(), 1);
316        batches.swap_remove(0)
317    }
318
319    /// Verifies the variant metadata is present in the schema for the specified
320    /// field name.
321    fn assert_variant_metadata(batch: &RecordBatch, field_name: &str) {
322        let schema = batch.schema();
323        let field = schema
324            .field_with_name(field_name)
325            .expect("could not find expected field");
326
327        // explicitly check the metadata so it is clear in the tests what the
328        // names are
329        let metadata_value = field
330            .metadata()
331            .get("ARROW:extension:name")
332            .expect("metadata does not exist");
333
334        assert_eq!(metadata_value, "arrow.parquet.variant");
335
336        // verify that `VariantType` also correctly finds the metadata
337        assert!(field.has_valid_extension_type::<VariantType>());
338    }
339
340    /// Read the specified test case filename from parquet-testing
341    /// See parquet-testing/shredded_variant/cases.json for more details
342    fn read_shredded_variant_test_case(name: &str) -> RecordBatch {
343        let case_file = PathBuf::from(parquet_test_data())
344            .join("..") // go up from data/ to parquet-testing/
345            .join("shredded_variant")
346            .join(name);
347        let case_file = std::fs::File::open(case_file).unwrap();
348        read_to_batch(case_file)
349    }
350}