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}