Skip to content

Commit 9d2193b

Browse files
committed
fix(storage): preserve Parquet schema annotations
1 parent b71dbbf commit 9d2193b

15 files changed

Lines changed: 730 additions & 80 deletions

File tree

src/query/catalog/src/plan/datasource/datasource_info/parquet.rs

Lines changed: 114 additions & 54 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,12 @@ use databend_common_storage::StageFileInfo;
2525
use databend_common_storage::StageFilesInfo;
2626
use databend_storages_common_table_meta::meta::ColumnStatistics;
2727
use parquet::arrow::ArrowSchemaConverter;
28+
use parquet::file::FOOTER_SIZE;
29+
use parquet::file::metadata::FileMetaData;
2830
use parquet::file::metadata::ParquetMetaData;
29-
use parquet::schema::parser::parse_message_type;
30-
use parquet::schema::printer::print_schema;
31+
use parquet::file::metadata::ParquetMetaDataReader;
32+
use parquet::file::metadata::ParquetMetaDataWriter;
3133
use parquet::schema::types::SchemaDescPtr;
32-
use parquet::schema::types::SchemaDescriptor;
3334
use serde::Deserialize;
3435
use serde::Serialize;
3536

@@ -59,6 +60,9 @@ pub struct ParquetTableInfo {
5960
pub table_info: TableInfo,
6061
pub arrow_schema: ArrowSchema,
6162
pub schema_descr: SchemaDescPtr,
63+
/// Whether `schema_descr` was rebuilt from `arrow_schema` because its serialized form
64+
/// could not be decoded. This is runtime diagnostic state and is not serialized.
65+
pub schema_descr_from_arrow_fallback: bool,
6266
pub files_to_read: Option<Vec<StageFileInfo>>,
6367
pub schema_from: String,
6468
pub compression_ratio: f64,
@@ -97,13 +101,15 @@ struct ParquetTableInfoSerde {
97101
impl Serialize for ParquetTableInfo {
98102
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
99103
where S: serde::Serializer {
104+
let schema_descr_bytes =
105+
schema_to_bytes(&self.schema_descr).map_err(serde::ser::Error::custom)?;
100106
ParquetTableInfoSerde {
101107
read_options: self.read_options,
102108
stage_info: self.stage_info.clone(),
103109
files_info: self.files_info.clone(),
104110
table_info: self.table_info.clone(),
105111
arrow_schema: self.arrow_schema.clone(),
106-
schema_descr_bytes: schema_to_bytes(&self.schema_descr),
112+
schema_descr_bytes,
107113
schema_descr_root: self.schema_descr.root_schema().name().to_string(),
108114
files_to_read: self.files_to_read.clone(),
109115
schema_from: self.schema_from.clone(),
@@ -118,12 +124,25 @@ impl<'de> Deserialize<'de> for ParquetTableInfo {
118124
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
119125
where D: serde::Deserializer<'de> {
120126
let helper = ParquetTableInfoSerde::deserialize(deserializer)?;
121-
let schema_descr = schema_from_bytes(&helper.schema_descr_bytes).or_else(|_| {
122-
ArrowSchemaConverter::new()
123-
.schema_root(&helper.schema_descr_root)
124-
.convert(&helper.arrow_schema)
125-
.map(Arc::new)
126-
});
127+
let (schema_descr, schema_descr_from_arrow_fallback) = match schema_from_bytes(
128+
&helper.schema_descr_bytes,
129+
) {
130+
Ok(schema_descr) => (Ok(schema_descr), false),
131+
Err(error) => {
132+
// Distributed query plans are expected to be exchanged by query nodes running
133+
// the same version. This fallback is best effort and may normalize annotations.
134+
log::warn!(
135+
"Failed to decode serialized Parquet schema metadata, falling back to the Arrow schema: {error}"
136+
);
137+
(
138+
ArrowSchemaConverter::new()
139+
.schema_root(&helper.schema_descr_root)
140+
.convert(&helper.arrow_schema)
141+
.map(Arc::new),
142+
true,
143+
)
144+
}
145+
};
127146
let schema_descr = schema_descr.map_err(|e| serde::de::Error::custom(e.to_string()))?;
128147

129148
Ok(Self {
@@ -133,6 +152,7 @@ impl<'de> Deserialize<'de> for ParquetTableInfo {
133152
table_info: helper.table_info,
134153
arrow_schema: helper.arrow_schema,
135154
schema_descr,
155+
schema_descr_from_arrow_fallback,
136156
files_to_read: helper.files_to_read,
137157
schema_from: helper.schema_from,
138158
compression_ratio: helper.compression_ratio,
@@ -145,16 +165,18 @@ impl<'de> Deserialize<'de> for ParquetTableInfo {
145165
}
146166

147167
fn schema_from_bytes(bytes: &[u8]) -> parquet::errors::Result<SchemaDescPtr> {
148-
let schema_string = String::from_utf8(bytes.to_vec())
149-
.map_err(|e| parquet::errors::ParquetError::General(e.to_string()))?;
150-
let schema = parse_message_type(&schema_string)?;
151-
Ok(Arc::new(SchemaDescriptor::new(Arc::new(schema))))
168+
ParquetMetaDataReader::decode_schema(bytes)
152169
}
153170

154-
fn schema_to_bytes(schema: &SchemaDescPtr) -> Vec<u8> {
171+
fn schema_to_bytes(schema: &SchemaDescPtr) -> parquet::errors::Result<Vec<u8>> {
172+
let file_metadata = FileMetaData::new(1, 0, None, None, schema.clone(), None);
173+
let metadata = ParquetMetaData::new(file_metadata, vec![]);
155174
let mut out = Vec::new();
156-
print_schema(&mut out, schema.root_schema());
157-
out
175+
ParquetMetaDataWriter::new(&mut out, &metadata).finish()?;
176+
// finish() appends the metadata length and PAR1 magic, while decode_schema()
177+
// expects only the Thrift FileMetaData payload.
178+
out.truncate(out.len() - FOOTER_SIZE);
179+
Ok(out)
158180
}
159181

160182
#[cfg(test)]
@@ -169,12 +191,13 @@ mod tests {
169191
use parquet::basic::Type as PhysicalType;
170192
use parquet::errors::ParquetError;
171193
use parquet::schema::parser::parse_message_type;
172-
use parquet::schema::printer::print_schema;
173194
use parquet::schema::types::SchemaDescPtr;
174195
use parquet::schema::types::SchemaDescriptor;
175196
use parquet::schema::types::Type;
176197

177198
use super::ParquetTableInfo;
199+
use super::schema_from_bytes;
200+
use super::schema_to_bytes;
178201

179202
fn make_desc() -> Result<SchemaDescPtr, ParquetError> {
180203
let mut fields = vec![];
@@ -232,11 +255,10 @@ mod tests {
232255
Ok(Arc::new(SchemaDescriptor::new(Arc::new(schema))))
233256
}
234257

235-
#[test]
236-
fn test_serde() {
237-
let schema_descr = make_desc().unwrap();
238-
let info = ParquetTableInfo {
239-
schema_descr: schema_descr.clone(),
258+
fn info_with(schema_descr: SchemaDescPtr, arrow_schema: ArrowSchema) -> ParquetTableInfo {
259+
ParquetTableInfo {
260+
schema_descr,
261+
schema_descr_from_arrow_fallback: false,
240262
read_options: Default::default(),
241263
stage_info: Default::default(),
242264
files_info: StageFilesInfo {
@@ -246,60 +268,98 @@ mod tests {
246268
},
247269
table_info: Default::default(),
248270
leaf_fields: Arc::new(vec![]),
249-
arrow_schema: ArrowSchema {
250-
fields: Default::default(),
251-
metadata: Default::default(),
252-
},
271+
arrow_schema,
253272
files_to_read: None,
254273
schema_from: "".to_string(),
255274
compression_ratio: 0.0,
256275
need_stats_provider: false,
257276
max_threads: 1,
258277
max_memory_usage: 10000,
259-
};
278+
}
279+
}
280+
281+
#[test]
282+
fn test_serde() {
283+
let schema_descr = make_desc().unwrap();
284+
let info = info_with(schema_descr.clone(), ArrowSchema::empty());
260285
let s = serde_json::to_string(&info).unwrap();
261286
let info = serde_json::from_str::<ParquetTableInfo>(&s).unwrap();
262287

263-
let mut original = Vec::new();
264-
print_schema(&mut original, schema_descr.root_schema());
265-
let mut roundtrip = Vec::new();
266-
print_schema(&mut roundtrip, info.schema_descr.root_schema());
288+
assert_eq!(schema_descr.root_schema(), info.schema_descr.root_schema());
289+
}
290+
291+
#[test]
292+
fn test_schema_bytes_roundtrip() {
293+
let schema_descr = make_desc().unwrap();
294+
let schema_bytes = schema_to_bytes(&schema_descr).unwrap();
295+
296+
assert!(!schema_bytes.ends_with(b"PAR1"));
297+
let decoded_schema = schema_from_bytes(&schema_bytes).unwrap();
298+
assert_eq!(schema_descr.root_schema(), decoded_schema.root_schema());
299+
}
300+
301+
#[test]
302+
fn test_serde_preserves_legacy_decimal_annotation() {
303+
let deal = Type::primitive_type_builder("deal", PhysicalType::FIXED_LEN_BYTE_ARRAY)
304+
.with_repetition(Repetition::OPTIONAL)
305+
.with_converted_type(ConvertedType::DECIMAL)
306+
.with_length(9)
307+
.with_precision(20)
308+
.with_scale(0)
309+
.build()
310+
.unwrap();
311+
let nested = Type::group_type_builder("nested")
312+
.with_repetition(Repetition::OPTIONAL)
313+
.with_fields(vec![Arc::new(deal)])
314+
.build()
315+
.unwrap();
316+
let schema = Type::group_type_builder("spark_schema")
317+
.with_fields(vec![Arc::new(nested)])
318+
.build()
319+
.unwrap();
320+
let schema_descr = Arc::new(SchemaDescriptor::new(Arc::new(schema)));
321+
assert!(!schema_to_bytes(&schema_descr).unwrap().ends_with(b"PAR1"));
322+
assert!(
323+
schema_descr
324+
.column(0)
325+
.self_type()
326+
.get_basic_info()
327+
.logical_type_ref()
328+
.is_none()
329+
);
330+
331+
let arrow_schema = parquet_to_arrow_schema(&schema_descr, None).unwrap();
332+
let info = info_with(schema_descr.clone(), arrow_schema);
333+
334+
let serialized = serde_json::to_string(&info).unwrap();
335+
let deserialized = serde_json::from_str::<ParquetTableInfo>(&serialized).unwrap();
267336

268337
assert_eq!(
269-
schema_descr.root_schema().name(),
270-
info.schema_descr.root_schema().name()
338+
schema_descr.root_schema(),
339+
deserialized.schema_descr.root_schema()
340+
);
341+
assert!(
342+
deserialized
343+
.schema_descr
344+
.column(0)
345+
.self_type()
346+
.get_basic_info()
347+
.logical_type_ref()
348+
.is_none()
271349
);
272-
assert_eq!(original, roundtrip)
273350
}
274351

275352
#[test]
276353
fn test_serde_falls_back_to_arrow_schema() {
277354
let schema_descr = make_arrow_compatible_desc().unwrap();
278355
let arrow_schema = parquet_to_arrow_schema(&schema_descr, None).unwrap();
279-
let info = ParquetTableInfo {
280-
schema_descr: schema_descr.clone(),
281-
read_options: Default::default(),
282-
stage_info: Default::default(),
283-
files_info: StageFilesInfo {
284-
path: "".to_string(),
285-
files: None,
286-
pattern: None,
287-
},
288-
table_info: Default::default(),
289-
leaf_fields: Arc::new(vec![]),
290-
arrow_schema,
291-
files_to_read: None,
292-
schema_from: "".to_string(),
293-
compression_ratio: 0.0,
294-
need_stats_provider: false,
295-
max_threads: 1,
296-
max_memory_usage: 10000,
297-
};
356+
let info = info_with(schema_descr.clone(), arrow_schema);
298357

299358
let mut json = serde_json::to_value(&info).unwrap();
300359
json["schema_descr_bytes"] = serde_json::json!(Vec::<u8>::from("invalid schema"));
301360

302361
let info = serde_json::from_value::<ParquetTableInfo>(json).unwrap();
362+
assert!(info.schema_descr_from_arrow_fallback);
303363

304364
assert_eq!(
305365
schema_descr.root_schema().name(),

src/query/storages/parquet/src/copy_into_table/reader.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,7 @@ impl RowGroupReaderForCopy {
119119
schema_descr,
120120
Some(arrow_schema),
121121
None,
122+
false,
122123
)
123124
.with_push_downs(Some(&pushdowns));
124125
reader_builder.build_output()?;

0 commit comments

Comments
 (0)