1use std::collections::HashMap;
19use std::sync::Arc;
20
21use arrow_array::RecordBatch;
22use arrow_array::builder::{
23 BooleanBuilder, GenericListBuilder, ListBuilder, PrimitiveBuilder, StringBuilder, StructBuilder,
24};
25use arrow_array::types::{Int32Type, Int64Type};
26use arrow_schema::{DataType, Field, Fields};
27use futures::{StreamExt, stream};
28
29use crate::Result;
30use crate::arrow::schema_to_arrow_schema;
31use crate::error::invalid_data;
32use crate::scan::ArrowRecordBatchStream;
33use crate::spec::{Datum, FieldSummary, ListType, NestedField, PrimitiveType, StructType, Type};
34use crate::table::Table;
35
36pub struct ManifestsTable<'a> {
38 table: &'a Table,
39}
40
41impl<'a> ManifestsTable<'a> {
42 pub fn new(table: &'a Table) -> Self {
44 Self { table }
45 }
46
47 pub fn schema(&self) -> crate::spec::Schema {
49 let fields = vec![
50 NestedField::new(14, "content", Type::Primitive(PrimitiveType::Int), true),
51 NestedField::new(1, "path", Type::Primitive(PrimitiveType::String), true),
52 NestedField::new(2, "length", Type::Primitive(PrimitiveType::Long), true),
53 NestedField::new(
54 3,
55 "partition_spec_id",
56 Type::Primitive(PrimitiveType::Int),
57 true,
58 ),
59 NestedField::new(
60 4,
61 "added_snapshot_id",
62 Type::Primitive(PrimitiveType::Long),
63 true,
64 ),
65 NestedField::new(
66 5,
67 "added_data_files_count",
68 Type::Primitive(PrimitiveType::Int),
69 true,
70 ),
71 NestedField::new(
72 6,
73 "existing_data_files_count",
74 Type::Primitive(PrimitiveType::Int),
75 true,
76 ),
77 NestedField::new(
78 7,
79 "deleted_data_files_count",
80 Type::Primitive(PrimitiveType::Int),
81 true,
82 ),
83 NestedField::new(
84 15,
85 "added_delete_files_count",
86 Type::Primitive(PrimitiveType::Int),
87 true,
88 ),
89 NestedField::new(
90 16,
91 "existing_delete_files_count",
92 Type::Primitive(PrimitiveType::Int),
93 true,
94 ),
95 NestedField::new(
96 17,
97 "deleted_delete_files_count",
98 Type::Primitive(PrimitiveType::Int),
99 true,
100 ),
101 NestedField::new(
102 8,
103 "partition_summaries",
104 Type::List(ListType {
105 element_field: Arc::new(NestedField::new(
106 9,
107 "item",
108 Type::Struct(StructType::new(vec![
109 Arc::new(NestedField::new(
110 10,
111 "contains_null",
112 Type::Primitive(PrimitiveType::Boolean),
113 true,
114 )),
115 Arc::new(NestedField::new(
116 11,
117 "contains_nan",
118 Type::Primitive(PrimitiveType::Boolean),
119 false,
120 )),
121 Arc::new(NestedField::new(
122 12,
123 "lower_bound",
124 Type::Primitive(PrimitiveType::String),
125 false,
126 )),
127 Arc::new(NestedField::new(
128 13,
129 "upper_bound",
130 Type::Primitive(PrimitiveType::String),
131 false,
132 )),
133 ])),
134 true,
135 )),
136 }),
137 true,
138 ),
139 ];
140
141 crate::spec::Schema::builder()
142 .with_fields(fields.into_iter().map(|f| f.into()))
143 .build()
144 .unwrap()
145 }
146
147 pub async fn scan(&self) -> Result<ArrowRecordBatchStream> {
149 let schema = schema_to_arrow_schema(&self.schema())?;
150
151 let mut content = PrimitiveBuilder::<Int32Type>::new();
152 let mut path = StringBuilder::new();
153 let mut length = PrimitiveBuilder::<Int64Type>::new();
154 let mut partition_spec_id = PrimitiveBuilder::<Int32Type>::new();
155 let mut added_snapshot_id = PrimitiveBuilder::<Int64Type>::new();
156 let mut added_data_files_count = PrimitiveBuilder::<Int32Type>::new();
157 let mut existing_data_files_count = PrimitiveBuilder::<Int32Type>::new();
158 let mut deleted_data_files_count = PrimitiveBuilder::<Int32Type>::new();
159 let mut added_delete_files_count = PrimitiveBuilder::<Int32Type>::new();
160 let mut existing_delete_files_count = PrimitiveBuilder::<Int32Type>::new();
161 let mut deleted_delete_files_count = PrimitiveBuilder::<Int32Type>::new();
162 let mut partition_summaries = self.partition_summary_builder()?;
163
164 if let Some(snapshot) = self.table.metadata().current_snapshot() {
165 let manifest_list = self.table.manifest_list_reader(snapshot).load().await?;
166 for manifest in manifest_list.entries() {
167 content.append_value(manifest.content as i32);
168 path.append_value(manifest.manifest_path.clone());
169 length.append_value(manifest.manifest_length);
170 partition_spec_id.append_value(manifest.partition_spec_id);
171 added_snapshot_id.append_value(manifest.added_snapshot_id);
172 added_data_files_count.append_value(manifest.added_files_count.unwrap_or(0) as i32);
173 existing_data_files_count
174 .append_value(manifest.existing_files_count.unwrap_or(0) as i32);
175 deleted_data_files_count
176 .append_value(manifest.deleted_files_count.unwrap_or(0) as i32);
177 added_delete_files_count
178 .append_value(manifest.added_files_count.unwrap_or(0) as i32);
179 existing_delete_files_count
180 .append_value(manifest.existing_files_count.unwrap_or(0) as i32);
181 deleted_delete_files_count
182 .append_value(manifest.deleted_files_count.unwrap_or(0) as i32);
183
184 let spec = self
185 .table
186 .metadata()
187 .partition_spec_by_id(manifest.partition_spec_id)
188 .ok_or_else(|| {
189 invalid_data!(
190 "Partition spec {} for manifest {} is not in table metadata",
191 manifest.partition_spec_id,
192 manifest.manifest_path
193 )
194 })?;
195 let spec_struct = spec.partition_type(self.table.metadata().current_schema())?;
196 self.append_partition_summaries(
197 &mut partition_summaries,
198 manifest.partitions.as_deref().unwrap_or(&[]),
199 spec_struct,
200 );
201 }
202 }
203
204 let batch = RecordBatch::try_new(Arc::new(schema), vec![
205 Arc::new(content.finish()),
206 Arc::new(path.finish()),
207 Arc::new(length.finish()),
208 Arc::new(partition_spec_id.finish()),
209 Arc::new(added_snapshot_id.finish()),
210 Arc::new(added_data_files_count.finish()),
211 Arc::new(existing_data_files_count.finish()),
212 Arc::new(deleted_data_files_count.finish()),
213 Arc::new(added_delete_files_count.finish()),
214 Arc::new(existing_delete_files_count.finish()),
215 Arc::new(deleted_delete_files_count.finish()),
216 Arc::new(partition_summaries.finish()),
217 ])?;
218 Ok(stream::iter(vec![Ok(batch)]).boxed())
219 }
220
221 fn partition_summary_builder(&self) -> Result<GenericListBuilder<i32, StructBuilder>> {
222 let schema = schema_to_arrow_schema(&self.schema())?;
223 let partition_summary_fields =
224 match schema.field_with_name("partition_summaries")?.data_type() {
225 DataType::List(list_type) => match list_type.data_type() {
226 DataType::Struct(fields) => fields.to_vec(),
227 _ => unreachable!(),
228 },
229 _ => unreachable!(),
230 };
231
232 let partition_summaries = ListBuilder::new(StructBuilder::from_fields(
233 Fields::from(partition_summary_fields.clone()),
234 0,
235 ))
236 .with_field(Arc::new(
237 Field::new_struct("item", partition_summary_fields, false).with_metadata(
238 HashMap::from([("PARQUET:field_id".to_string(), "9".to_string())]),
239 ),
240 ));
241
242 Ok(partition_summaries)
243 }
244
245 fn append_partition_summaries(
246 &self,
247 builder: &mut GenericListBuilder<i32, StructBuilder>,
248 partitions: &[FieldSummary],
249 partition_struct: StructType,
250 ) {
251 let partition_summaries_builder = builder.values();
252 for (summary, field) in partitions.iter().zip(partition_struct.fields()) {
253 partition_summaries_builder
254 .field_builder::<BooleanBuilder>(0)
255 .unwrap()
256 .append_value(summary.contains_null);
257 partition_summaries_builder
258 .field_builder::<BooleanBuilder>(1)
259 .unwrap()
260 .append_option(summary.contains_nan);
261
262 let field_type = field.field_type.as_primitive_type().unwrap();
263 for (index, bound) in [(2, &summary.lower_bound), (3, &summary.upper_bound)] {
264 let bound = bound
266 .as_ref()
267 .filter(|_| *field_type != PrimitiveType::Unknown)
268 .map(|v| {
269 Datum::try_from_bytes(v, field_type.clone())
270 .unwrap()
271 .to_string()
272 });
273 partition_summaries_builder
274 .field_builder::<StringBuilder>(index)
275 .unwrap()
276 .append_option(bound);
277 }
278 partition_summaries_builder.append(true);
279 }
280 builder.append(true);
281 }
282}
283
284#[cfg(test)]
285mod tests {
286 use std::sync::Arc;
287
288 use arrow_array::{Array, ListArray, StructArray};
289 use expect_test::expect;
290 use futures::TryStreamExt;
291
292 use crate::spec::TableMetadata;
293 use crate::test_utils::check_record_batches;
294 use crate::test_utils::scan::TableTestFixture;
295
296 #[tokio::test]
297 async fn test_manifests_table() {
298 let mut fixture = TableTestFixture::new();
299 fixture.setup_manifest_files().await;
300
301 let record_batch = fixture.table.inspect().manifests().scan().await.unwrap();
302
303 check_record_batches(
304 record_batch.try_collect::<Vec<_>>().await.unwrap(),
305 expect![[r#"
306 Field { "content": Int32, metadata: {"PARQUET:field_id": "14"} },
307 Field { "path": Utf8, metadata: {"PARQUET:field_id": "1"} },
308 Field { "length": Int64, metadata: {"PARQUET:field_id": "2"} },
309 Field { "partition_spec_id": Int32, metadata: {"PARQUET:field_id": "3"} },
310 Field { "added_snapshot_id": Int64, metadata: {"PARQUET:field_id": "4"} },
311 Field { "added_data_files_count": Int32, metadata: {"PARQUET:field_id": "5"} },
312 Field { "existing_data_files_count": Int32, metadata: {"PARQUET:field_id": "6"} },
313 Field { "deleted_data_files_count": Int32, metadata: {"PARQUET:field_id": "7"} },
314 Field { "added_delete_files_count": Int32, metadata: {"PARQUET:field_id": "15"} },
315 Field { "existing_delete_files_count": Int32, metadata: {"PARQUET:field_id": "16"} },
316 Field { "deleted_delete_files_count": Int32, metadata: {"PARQUET:field_id": "17"} },
317 Field { "partition_summaries": List(non-null Struct("contains_null": non-null Boolean, metadata: {"PARQUET:field_id": "10"}, "contains_nan": Boolean, metadata: {"PARQUET:field_id": "11"}, "lower_bound": Utf8, metadata: {"PARQUET:field_id": "12"}, "upper_bound": Utf8, metadata: {"PARQUET:field_id": "13"}), metadata: {"PARQUET:field_id": "9"}), metadata: {"PARQUET:field_id": "8"} }"#]],
318 expect![[r#"
319 content: PrimitiveArray<Int32>
320 [
321 0,
322 ],
323 path: (skipped),
324 length: (skipped),
325 partition_spec_id: PrimitiveArray<Int32>
326 [
327 0,
328 ],
329 added_snapshot_id: PrimitiveArray<Int64>
330 [
331 3055729675574597004,
332 ],
333 added_data_files_count: PrimitiveArray<Int32>
334 [
335 1,
336 ],
337 existing_data_files_count: PrimitiveArray<Int32>
338 [
339 1,
340 ],
341 deleted_data_files_count: PrimitiveArray<Int32>
342 [
343 1,
344 ],
345 added_delete_files_count: PrimitiveArray<Int32>
346 [
347 1,
348 ],
349 existing_delete_files_count: PrimitiveArray<Int32>
350 [
351 1,
352 ],
353 deleted_delete_files_count: PrimitiveArray<Int32>
354 [
355 1,
356 ],
357 partition_summaries: ListArray
358 [
359 StructArray
360 -- validity:
361 [
362 valid,
363 ]
364 [
365 -- child 0: "contains_null" (Boolean)
366 BooleanArray
367 [
368 false,
369 ]
370 -- child 1: "contains_nan" (Boolean)
371 BooleanArray
372 [
373 false,
374 ]
375 -- child 2: "lower_bound" (Utf8)
376 StringArray
377 [
378 "100",
379 ]
380 -- child 3: "upper_bound" (Utf8)
381 StringArray
382 [
383 "300",
384 ]
385 ],
386 ]"#]],
387 &["path", "length"],
388 Some("path"),
389 );
390 }
391
392 #[tokio::test]
393 async fn test_manifests_table_with_dropped_partition_source_column() {
394 let mut fixture = TableTestFixture::new();
395 fixture.setup_manifest_files().await;
396
397 let mut metadata = serde_json::to_value(fixture.table.metadata()).unwrap();
401 let current_schema_id = metadata["current-schema-id"].clone();
402 let schemas = metadata["schemas"].as_array_mut().unwrap();
403 for schema in schemas {
404 if schema["schema-id"] == current_schema_id {
405 let fields = schema["fields"].as_array_mut().unwrap();
406 fields.retain(|field| field["id"] != 1);
407 let identifier_ids = schema["identifier-field-ids"].as_array_mut().unwrap();
408 identifier_ids.retain(|id| *id != 1);
409 }
410 }
411
412 metadata["partition-specs"]
413 .as_array_mut()
414 .unwrap()
415 .push(serde_json::json!({"spec-id": 1, "fields": []}));
416 metadata["default-spec-id"] = serde_json::json!(1);
417
418 let metadata: TableMetadata = serde_json::from_value(metadata).unwrap();
419 let table = fixture.table.clone().with_metadata(Arc::new(metadata));
420
421 let batches: Vec<_> = table
422 .inspect()
423 .manifests()
424 .scan()
425 .await
426 .unwrap()
427 .try_collect()
428 .await
429 .unwrap();
430 assert_eq!(batches.len(), 1);
431 assert_eq!(batches[0].num_rows(), 1);
432 let summaries = batches[0]
433 .column_by_name("partition_summaries")
434 .unwrap()
435 .as_any()
436 .downcast_ref::<ListArray>()
437 .unwrap();
438 let summary = summaries.value(0);
439 let summary = summary.as_any().downcast_ref::<StructArray>().unwrap();
440 assert_eq!(summary.len(), 1);
441 assert!(summary.column_by_name("contains_null").unwrap().is_valid(0));
442 assert!(summary.column_by_name("contains_nan").unwrap().is_valid(0));
443 assert!(summary.column_by_name("lower_bound").unwrap().is_null(0));
444 assert!(summary.column_by_name("upper_bound").unwrap().is_null(0));
445 }
446}