iceberg/spec/manifest_list/
mod.rs1mod _const_schema;
21pub(super) mod _serde;
22mod manifest_file;
23mod reader;
24mod writer;
25
26use apache_avro::types::Value;
27use apache_avro::{Reader, from_value};
28pub use manifest_file::*;
29pub use reader::*;
30pub use serde_bytes::ByteBuf;
31pub use writer::*;
32
33use self::_const_schema::MANIFEST_LIST_AVRO_SCHEMA_V1;
34use super::FormatVersion;
35use crate::error::Result;
36
37pub const UNASSIGNED_SEQUENCE_NUMBER: i64 = -1;
39
40#[derive(Debug, Clone, PartialEq)]
54pub struct ManifestList {
55 entries: Vec<ManifestFile>,
57}
58
59impl ManifestList {
60 pub fn parse_with_version(bs: &[u8], version: FormatVersion) -> Result<ManifestList> {
62 match version {
63 FormatVersion::V1 => {
64 let reader = Reader::builder(bs)
65 .reader_schema(&MANIFEST_LIST_AVRO_SCHEMA_V1)
66 .build()?;
67 let values = Value::Array(reader.collect::<std::result::Result<Vec<Value>, _>>()?);
68 from_value::<_serde::ManifestListV1>(&values)?.try_into()
69 }
70 FormatVersion::V2 => {
71 let reader = Reader::new(bs)?;
72 let values = Value::Array(reader.collect::<std::result::Result<Vec<Value>, _>>()?);
73 from_value::<_serde::ManifestListV2>(&values)?.try_into()
74 }
75 FormatVersion::V3 => {
76 let reader = Reader::new(bs)?;
77 let values = Value::Array(reader.collect::<std::result::Result<Vec<Value>, _>>()?);
78 from_value::<_serde::ManifestListV3>(&values)?.try_into()
79 }
80 }
81 }
82
83 pub fn entries(&self) -> &[ManifestFile] {
85 &self.entries
86 }
87
88 pub fn consume_entries(self) -> impl IntoIterator<Item = ManifestFile> {
90 Box::new(self.entries.into_iter())
91 }
92}
93
94#[cfg(test)]
95mod test {
96 use std::fs;
97
98 use apache_avro::{Codec, Writer};
99 use tempfile::TempDir;
100
101 use super::_const_schema::MANIFEST_LIST_AVRO_SCHEMA_V2;
102 use super::_serde::ManifestFileV2;
103 use super::*;
104 use crate::io::FileIO;
105 use crate::spec::{Datum, FieldSummary, ManifestContentType, ManifestFile};
106
107 #[tokio::test]
108 async fn test_parse_manifest_list_v1() {
109 let manifest_list = ManifestList {
110 entries: vec![
111 ManifestFile {
112 manifest_path: "/opt/bitnami/spark/warehouse/db/table/metadata/10d28031-9739-484c-92db-cdf2975cead4-m0.avro".to_string(),
113 manifest_length: 5806,
114 partition_spec_id: 0,
115 content: ManifestContentType::Data,
116 sequence_number: 0,
117 min_sequence_number: 0,
118 added_snapshot_id: 1646658105718557341,
119 added_files_count: Some(3),
120 existing_files_count: Some(0),
121 deleted_files_count: Some(0),
122 added_rows_count: Some(3),
123 existing_rows_count: Some(0),
124 deleted_rows_count: Some(0),
125 partitions: Some(vec![]),
126 key_metadata: None,
127 first_row_id: None,
128 }
129 ]
130 };
131
132 let file_io = FileIO::new_with_fs();
133
134 let tmp_dir = TempDir::new().unwrap();
135 let file_name = "simple_manifest_list_v1.avro";
136 let full_path = format!("{}/{}", tmp_dir.path().to_str().unwrap(), file_name);
137
138 let mut writer = ManifestListWriter::v1(
139 file_io
140 .new_output(full_path.clone())
141 .unwrap()
142 .writer()
143 .await
144 .unwrap(),
145 1646658105718557341,
146 Some(1646658105718557341),
147 );
148
149 writer
150 .add_manifests(manifest_list.entries.clone().into_iter())
151 .unwrap();
152 writer.close().await.unwrap();
153
154 let bs = fs::read(full_path).expect("read_file must succeed");
155
156 let parsed_manifest_list =
157 ManifestList::parse_with_version(&bs, FormatVersion::V1).unwrap();
158
159 assert_eq!(manifest_list, parsed_manifest_list);
160 }
161
162 #[tokio::test]
163 async fn test_parse_manifest_list_v2() {
164 let manifest_list = ManifestList {
165 entries: vec![
166 ManifestFile {
167 manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m0.avro".to_string(),
168 manifest_length: 6926,
169 partition_spec_id: 1,
170 content: ManifestContentType::Data,
171 sequence_number: 1,
172 min_sequence_number: 1,
173 added_snapshot_id: 377075049360453639,
174 added_files_count: Some(1),
175 existing_files_count: Some(0),
176 deleted_files_count: Some(0),
177 added_rows_count: Some(3),
178 existing_rows_count: Some(0),
179 deleted_rows_count: Some(0),
180 partitions: Some(
181 vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
182 ),
183 key_metadata: None,
184 first_row_id: None,
185 },
186 ManifestFile {
187 manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m1.avro".to_string(),
188 manifest_length: 6926,
189 partition_spec_id: 2,
190 content: ManifestContentType::Data,
191 sequence_number: 1,
192 min_sequence_number: 1,
193 added_snapshot_id: 377075049360453639,
194 added_files_count: Some(1),
195 existing_files_count: Some(0),
196 deleted_files_count: Some(0),
197 added_rows_count: Some(3),
198 existing_rows_count: Some(0),
199 deleted_rows_count: Some(0),
200 partitions: Some(
201 vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::float(1.1_f32).to_bytes().unwrap()), upper_bound: Some(Datum::float(2.1_f32).to_bytes().unwrap())}]
202 ),
203 key_metadata: None,
204 first_row_id: None,
205 }
206 ]
207 };
208
209 let file_io = FileIO::new_with_fs();
210
211 let tmp_dir = TempDir::new().unwrap();
212 let file_name = "simple_manifest_list_v1.avro";
213 let full_path = format!("{}/{}", tmp_dir.path().to_str().unwrap(), file_name);
214
215 let mut writer = ManifestListWriter::v2(
216 file_io
217 .new_output(full_path.clone())
218 .unwrap()
219 .writer()
220 .await
221 .unwrap(),
222 1646658105718557341,
223 Some(1646658105718557341),
224 1,
225 );
226
227 writer
228 .add_manifests(manifest_list.entries.clone().into_iter())
229 .unwrap();
230 writer.close().await.unwrap();
231
232 let bs = fs::read(full_path).expect("read_file must succeed");
233
234 let parsed_manifest_list =
235 ManifestList::parse_with_version(&bs, FormatVersion::V2).unwrap();
236
237 assert_eq!(manifest_list, parsed_manifest_list);
238 }
239
240 #[test]
241 fn test_parse_snappy_manifest_list_v2() {
242 let manifest_list = ManifestList {
243 entries: vec![ManifestFile {
244 manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/snappy-m0.avro".to_string(),
245 manifest_length: 6926,
246 partition_spec_id: 1,
247 content: ManifestContentType::Data,
248 sequence_number: 1,
249 min_sequence_number: 1,
250 added_snapshot_id: 377075049360453639,
251 added_files_count: Some(1),
252 existing_files_count: Some(0),
253 deleted_files_count: Some(0),
254 added_rows_count: Some(3),
255 existing_rows_count: Some(0),
256 deleted_rows_count: Some(0),
257 partitions: Some(vec![FieldSummary {
258 contains_null: false,
259 contains_nan: Some(false),
260 lower_bound: Some(Datum::long(1).to_bytes().unwrap()),
261 upper_bound: Some(Datum::long(1).to_bytes().unwrap()),
262 }]),
263 key_metadata: None,
264 first_row_id: None,
265 }],
266 };
267
268 let manifest_entry: ManifestFileV2 = manifest_list.entries[0].clone().try_into().unwrap();
269 let mut writer =
270 Writer::with_codec(&MANIFEST_LIST_AVRO_SCHEMA_V2, Vec::new(), Codec::Snappy).unwrap();
271 writer.append_ser(manifest_entry).unwrap();
272 let bs = writer.into_inner().unwrap();
273
274 let parsed_manifest_list =
275 ManifestList::parse_with_version(&bs, FormatVersion::V2).unwrap();
276
277 assert_eq!(manifest_list, parsed_manifest_list);
278 }
279
280 #[tokio::test]
281 async fn test_parse_manifest_list_v3() {
282 let manifest_list = ManifestList {
283 entries: vec![
284 ManifestFile {
285 manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m0.avro".to_string(),
286 manifest_length: 6926,
287 partition_spec_id: 1,
288 content: ManifestContentType::Data,
289 sequence_number: 1,
290 min_sequence_number: 1,
291 added_snapshot_id: 377075049360453639,
292 added_files_count: Some(1),
293 existing_files_count: Some(0),
294 deleted_files_count: Some(0),
295 added_rows_count: Some(3),
296 existing_rows_count: Some(0),
297 deleted_rows_count: Some(0),
298 partitions: Some(
299 vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
300 ),
301 key_metadata: None,
302 first_row_id: Some(10),
303 },
304 ManifestFile {
305 manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m1.avro".to_string(),
306 manifest_length: 6926,
307 partition_spec_id: 2,
308 content: ManifestContentType::Data,
309 sequence_number: 1,
310 min_sequence_number: 1,
311 added_snapshot_id: 377075049360453639,
312 added_files_count: Some(1),
313 existing_files_count: Some(0),
314 deleted_files_count: Some(0),
315 added_rows_count: Some(3),
316 existing_rows_count: Some(0),
317 deleted_rows_count: Some(0),
318 partitions: Some(
319 vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::float(1.1_f32).to_bytes().unwrap()), upper_bound: Some(Datum::float(2.1_f32).to_bytes().unwrap())}]
320 ),
321 key_metadata: None,
322 first_row_id: Some(13),
323 }
324 ]
325 };
326
327 let file_io = FileIO::new_with_fs();
328
329 let tmp_dir = TempDir::new().unwrap();
330 let file_name = "simple_manifest_list_v3.avro";
331 let full_path = format!("{}/{}", tmp_dir.path().to_str().unwrap(), file_name);
332
333 let mut writer = ManifestListWriter::v3(
334 file_io
335 .new_output(full_path.clone())
336 .unwrap()
337 .writer()
338 .await
339 .unwrap(),
340 377075049360453639,
341 Some(377075049360453639),
342 1,
343 Some(10),
344 );
345
346 writer
347 .add_manifests(manifest_list.entries.clone().into_iter())
348 .unwrap();
349 writer.close().await.unwrap();
350
351 let bs = fs::read(full_path).expect("read_file must succeed");
352
353 let parsed_manifest_list =
354 ManifestList::parse_with_version(&bs, FormatVersion::V3).unwrap();
355
356 assert_eq!(manifest_list, parsed_manifest_list);
357 }
358}