Skip to main content

iceberg/spec/manifest_list/
mod.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//! ManifestList for Iceberg.
19
20mod _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
37/// Placeholder for sequence number. The field with this value must be replaced with the actual sequence number before it write.
38pub const UNASSIGNED_SEQUENCE_NUMBER: i64 = -1;
39
40/// Snapshots are embedded in table metadata, but the list of manifests for a
41/// snapshot are stored in a separate manifest list file.
42///
43/// A new manifest list is written for each attempt to commit a snapshot
44/// because the list of manifests always changes to produce a new snapshot.
45/// When a manifest list is written, the (optimistic) sequence number of the
46/// snapshot is written for all new manifest files tracked by the list.
47///
48/// A manifest list includes summary metadata that can be used to avoid
49/// scanning all of the manifests in a snapshot when planning a table scan.
50/// This includes the number of added, existing, and deleted files, and a
51/// summary of values for each field of the partition spec used to write the
52/// manifest.
53#[derive(Debug, Clone, PartialEq)]
54pub struct ManifestList {
55    /// Entries in a manifest list.
56    entries: Vec<ManifestFile>,
57}
58
59impl ManifestList {
60    /// Parse manifest list from bytes.
61    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    /// Get the entries in the manifest list.
84    pub fn entries(&self) -> &[ManifestFile] {
85        &self.entries
86    }
87
88    /// Take ownership of the entries in the manifest list, consuming it
89    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}