1use crate::file::metadata::thrift::FileMeta;
19use crate::file::metadata::{ColumnChunkMetaData, PageIndex, RowGroupMetaData};
20use crate::schema::types::{SchemaDescPtr, SchemaDescriptor};
21use crate::{
22 basic::ColumnOrder,
23 file::metadata::{FileMetaData, ParquetMetaDataBuilder},
24};
25#[cfg(feature = "encryption")]
26use crate::{
27 encryption::{
28 encrypt::{FileEncryptor, encrypt_thrift_object, write_signed_plaintext_thrift_object},
29 modules::{ModuleType, create_footer_aad, create_module_aad},
30 },
31 file::column_crypto_metadata::ColumnCryptoMetaData,
32 file::metadata::thrift::encryption::{AesGcmV1, EncryptionAlgorithm, FileCryptoMetaData},
33};
34use crate::{errors::Result, file::page_index::column_index::ColumnIndexMetaData};
35
36use crate::{
37 file::writer::{TrackedWrite, get_file_magic},
38 parquet_thrift::WriteThrift,
39};
40use crate::{
41 file::{
42 metadata::{KeyValue, ParquetMetaData},
43 page_index::offset_index::OffsetIndexMetaData,
44 },
45 parquet_thrift::ThriftCompactOutputProtocol,
46};
47use std::io::Write;
48use std::sync::Arc;
49
50pub(crate) struct ThriftMetadataWriter<'a, W: Write> {
54 buf: &'a mut TrackedWrite<W>,
55 schema_descr: &'a SchemaDescPtr,
56 row_groups: Vec<RowGroupMetaData>,
57 column_indexes: Option<Vec<Vec<Option<ColumnIndexMetaData>>>>,
58 offset_indexes: Option<Vec<Vec<Option<OffsetIndexMetaData>>>>,
59 key_value_metadata: Option<Vec<KeyValue>>,
60 created_by: Option<String>,
61 object_writer: MetadataObjectWriter,
62 writer_version: i32,
63 write_path_in_schema: bool,
64}
65
66impl<'a, W: Write> ThriftMetadataWriter<'a, W> {
67 fn write_offset_indexes(
73 &mut self,
74 offset_indexes: &[Vec<Option<OffsetIndexMetaData>>],
75 ) -> Result<()> {
76 for (row_group_idx, row_group) in self.row_groups.iter_mut().enumerate() {
80 for (column_idx, column_metadata) in row_group.columns.iter_mut().enumerate() {
81 if let Some(offset_index) = &offset_indexes[row_group_idx][column_idx] {
82 let start_offset = self.buf.bytes_written();
83 self.object_writer.write_offset_index(
84 offset_index,
85 column_metadata,
86 row_group_idx,
87 column_idx,
88 &mut self.buf,
89 )?;
90 let end_offset = self.buf.bytes_written();
91 column_metadata.offset_index_offset = Some(start_offset as i64);
93 column_metadata.offset_index_length = Some((end_offset - start_offset) as i32);
94 }
95 }
96 }
97 Ok(())
98 }
99
100 fn write_column_indexes(
106 &mut self,
107 column_indexes: &[Vec<Option<ColumnIndexMetaData>>],
108 ) -> Result<()> {
109 for (row_group_idx, row_group) in self.row_groups.iter_mut().enumerate() {
113 for (column_idx, column_metadata) in row_group.columns.iter_mut().enumerate() {
114 if let Some(column_index) = &column_indexes[row_group_idx][column_idx] {
115 let start_offset = self.buf.bytes_written();
116 if self.object_writer.write_column_index(
118 column_index,
119 column_metadata,
120 row_group_idx,
121 column_idx,
122 &mut self.buf,
123 )? {
124 let end_offset = self.buf.bytes_written();
125 column_metadata.column_index_offset = Some(start_offset as i64);
127 column_metadata.column_index_length =
128 Some((end_offset - start_offset) as i32);
129 }
130 }
131 }
132 }
133 Ok(())
134 }
135
136 fn finalize_column_indexes(&mut self) -> Result<Option<Vec<Vec<Option<ColumnIndexMetaData>>>>> {
138 let column_indexes = std::mem::take(&mut self.column_indexes);
139
140 if let Some(column_indexes) = column_indexes.as_ref() {
142 self.write_column_indexes(column_indexes)?;
143 }
144
145 let all_none = column_indexes
147 .as_ref()
148 .is_some_and(|ci| ci.iter().all(|cii| cii.iter().all(|idx| idx.is_none())));
149
150 if all_none {
151 Ok(None)
152 } else {
153 Ok(column_indexes)
154 }
155 }
156
157 fn finalize_offset_indexes(&mut self) -> Result<Option<Vec<Vec<Option<OffsetIndexMetaData>>>>> {
159 let offset_indexes = std::mem::take(&mut self.offset_indexes);
160
161 if let Some(offset_indexes) = offset_indexes.as_ref() {
163 self.write_offset_indexes(offset_indexes)?;
164 }
165
166 let all_none = offset_indexes
168 .as_ref()
169 .is_some_and(|oi| oi.iter().all(|oii| oii.iter().all(|idx| idx.is_none())));
170
171 if all_none {
172 Ok(None)
173 } else {
174 Ok(offset_indexes)
175 }
176 }
177
178 pub fn finish(mut self) -> Result<ParquetMetaData> {
180 let num_rows = self.row_groups.iter().map(|x| x.num_rows).sum();
181
182 let column_indexes = self.finalize_column_indexes()?;
184 let offset_indexes = self.finalize_offset_indexes()?;
185
186 let column_orders = self
188 .schema_descr
189 .columns()
190 .iter()
191 .map(|col| {
192 ColumnOrder::column_order_for_type(
193 col.logical_type_ref(),
194 col.converted_type(),
195 col.physical_type(),
196 )
197 })
198 .collect();
199
200 let column_orders = Some(column_orders);
204
205 let (row_groups, unencrypted_row_groups) = self
206 .object_writer
207 .apply_row_group_encryption(self.row_groups)?;
208
209 #[cfg(feature = "encryption")]
210 let (encryption_algorithm, footer_signing_key_metadata) =
211 self.object_writer.get_plaintext_footer_crypto_metadata();
212 #[cfg(feature = "encryption")]
213 let file_metadata = FileMetaData::new(
214 self.writer_version,
215 num_rows,
216 self.created_by,
217 self.key_value_metadata,
218 self.schema_descr.clone(),
219 column_orders,
220 )
221 .with_encryption_algorithm(encryption_algorithm)
222 .with_footer_signing_key_metadata(footer_signing_key_metadata);
223
224 #[cfg(not(feature = "encryption"))]
225 let file_metadata = FileMetaData::new(
226 self.writer_version,
227 num_rows,
228 self.created_by,
229 self.key_value_metadata,
230 self.schema_descr.clone(),
231 column_orders,
232 );
233
234 let file_meta = FileMeta {
235 file_metadata: &file_metadata,
236 row_groups: &row_groups,
237 write_path_in_schema: self.write_path_in_schema,
238 };
239
240 let start_pos = self.buf.bytes_written();
242 self.object_writer
243 .write_file_metadata(&file_meta, &mut self.buf)?;
244 let end_pos = self.buf.bytes_written();
245
246 let metadata_len = (end_pos - start_pos) as u32;
248
249 self.buf.write_all(&metadata_len.to_le_bytes())?;
250 self.buf.write_all(self.object_writer.get_file_magic())?;
251
252 let builder = ParquetMetaDataBuilder::new(file_metadata).set_page_index(Some(Arc::new(
257 PageIndex::new(column_indexes, offset_indexes),
258 )));
259
260 Ok(match unencrypted_row_groups {
261 Some(rg) => builder.set_row_groups(rg).build(),
262 None => builder.set_row_groups(row_groups).build(),
263 })
264 }
265
266 pub fn new(
267 buf: &'a mut TrackedWrite<W>,
268 schema_descr: &'a SchemaDescPtr,
269 row_groups: Vec<RowGroupMetaData>,
270 created_by: Option<String>,
271 writer_version: i32,
272 write_path_in_schema: bool,
273 ) -> Self {
274 Self {
275 buf,
276 schema_descr,
277 row_groups,
278 column_indexes: None,
279 offset_indexes: None,
280 key_value_metadata: None,
281 created_by,
282 object_writer: Default::default(),
283 writer_version,
284 write_path_in_schema,
285 }
286 }
287
288 pub fn with_column_indexes(
289 mut self,
290 column_indexes: Vec<Vec<Option<ColumnIndexMetaData>>>,
291 ) -> Self {
292 self.column_indexes = Some(column_indexes);
293 self
294 }
295
296 pub fn with_offset_indexes(
297 mut self,
298 offset_indexes: Vec<Vec<Option<OffsetIndexMetaData>>>,
299 ) -> Self {
300 self.offset_indexes = Some(offset_indexes);
301 self
302 }
303
304 pub fn with_key_value_metadata(mut self, key_value_metadata: Vec<KeyValue>) -> Self {
305 self.key_value_metadata = Some(key_value_metadata);
306 self
307 }
308
309 #[cfg(feature = "encryption")]
310 pub fn with_file_encryptor(mut self, file_encryptor: Option<Arc<FileEncryptor>>) -> Self {
311 self.object_writer = self.object_writer.with_file_encryptor(file_encryptor);
312 self
313 }
314}
315
316pub struct ParquetMetaDataWriter<'a, W: Write> {
406 buf: TrackedWrite<W>,
407 metadata: &'a ParquetMetaData,
408 write_path_in_schema: bool,
409}
410
411impl<'a, W: Write> ParquetMetaDataWriter<'a, W> {
412 pub fn new(buf: W, metadata: &'a ParquetMetaData) -> Self {
420 Self::new_with_tracked(TrackedWrite::new(buf), metadata)
421 }
422
423 pub fn new_with_tracked(buf: TrackedWrite<W>, metadata: &'a ParquetMetaData) -> Self {
430 Self {
431 buf,
432 metadata,
433 write_path_in_schema: true,
434 }
435 }
436
437 pub fn with_write_path_in_schema(self, val: bool) -> Self {
440 Self {
441 write_path_in_schema: val,
442 ..self
443 }
444 }
445
446 pub fn finish(mut self) -> Result<()> {
448 let file_metadata = self.metadata.file_metadata();
449
450 let schema = Arc::new(file_metadata.schema().clone());
451 let schema_descr = Arc::new(SchemaDescriptor::new(schema.clone()));
452 let created_by = file_metadata.created_by().map(str::to_string);
453
454 let row_groups = self.metadata.row_groups.clone();
455
456 let key_value_metadata = file_metadata.key_value_metadata().cloned();
457
458 let mut encoder = ThriftMetadataWriter::new(
459 &mut self.buf,
460 &schema_descr,
461 row_groups,
462 created_by,
463 file_metadata.version(),
464 self.write_path_in_schema,
465 );
466
467 if let Some(page_index_arc) = self.metadata.page_index.as_ref()
471 && let Some(page_index) = page_index_arc
472 .as_any()
473 .downcast_ref::<crate::file::metadata::PageIndex>()
474 {
475 if let Some(column_indexes) = page_index.column_indexes_raw() {
476 encoder = encoder.with_column_indexes(column_indexes.clone());
477 }
478
479 if let Some(offset_indexes) = page_index.offset_indexes_raw() {
480 encoder = encoder.with_offset_indexes(offset_indexes.clone());
481 }
482 }
483
484 if let Some(key_value_metadata) = key_value_metadata {
485 encoder = encoder.with_key_value_metadata(key_value_metadata);
486 }
487 encoder.finish()?;
488
489 Ok(())
490 }
491}
492
493#[derive(Debug, Default)]
494struct MetadataObjectWriter {
495 #[cfg(feature = "encryption")]
496 file_encryptor: Option<Arc<FileEncryptor>>,
497}
498
499impl MetadataObjectWriter {
500 #[inline]
501 fn write_thrift_object(object: &impl WriteThrift, sink: impl Write) -> Result<()> {
502 let mut protocol = ThriftCompactOutputProtocol::new(sink);
503 object.write_thrift(&mut protocol)?;
504 Ok(())
505 }
506}
507
508#[cfg(not(feature = "encryption"))]
510impl MetadataObjectWriter {
511 fn write_file_metadata(&self, file_metadata: &FileMeta, sink: impl Write) -> Result<()> {
515 Self::write_thrift_object(file_metadata, sink)
516 }
517
518 fn write_offset_index(
522 &self,
523 offset_index: &OffsetIndexMetaData,
524 _column_chunk: &ColumnChunkMetaData,
525 _row_group_idx: usize,
526 _column_idx: usize,
527 sink: impl Write,
528 ) -> Result<()> {
529 Self::write_thrift_object(offset_index, sink)
530 }
531
532 fn write_column_index(
538 &self,
539 column_index: &ColumnIndexMetaData,
540 _column_chunk: &ColumnChunkMetaData,
541 _row_group_idx: usize,
542 _column_idx: usize,
543 sink: impl Write,
544 ) -> Result<bool> {
545 Self::write_thrift_object(column_index, sink)?;
546 Ok(true)
547 }
548
549 fn apply_row_group_encryption(
551 &self,
552 row_groups: Vec<RowGroupMetaData>,
553 ) -> Result<(Vec<RowGroupMetaData>, Option<Vec<RowGroupMetaData>>)> {
554 Ok((row_groups, None))
555 }
556
557 pub fn get_file_magic(&self) -> &[u8; 4] {
559 get_file_magic()
560 }
561}
562
563#[cfg(feature = "encryption")]
565impl MetadataObjectWriter {
566 fn with_file_encryptor(mut self, encryptor: Option<Arc<FileEncryptor>>) -> Self {
568 self.file_encryptor = encryptor;
569 self
570 }
571
572 fn write_file_metadata(&self, file_metadata: &FileMeta, mut sink: impl Write) -> Result<()> {
576 match self.file_encryptor.as_ref() {
577 Some(file_encryptor) if file_encryptor.properties().encrypt_footer() => {
578 let crypto_metadata = Self::file_crypto_metadata(file_encryptor)?;
580 let mut protocol = ThriftCompactOutputProtocol::new(&mut sink);
581 crypto_metadata.write_thrift(&mut protocol)?;
582
583 let aad = create_footer_aad(file_encryptor.file_aad())?;
585 let mut encryptor = file_encryptor.get_footer_encryptor()?;
586 encrypt_thrift_object(file_metadata, &mut encryptor, &mut sink, &aad)
587 }
588 Some(file_encryptor) if file_metadata.file_metadata.encryption_algorithm.is_some() => {
589 let aad = create_footer_aad(file_encryptor.file_aad())?;
590 let mut encryptor = file_encryptor.get_footer_encryptor()?;
591 write_signed_plaintext_thrift_object(file_metadata, &mut encryptor, &mut sink, &aad)
592 }
593 _ => Self::write_thrift_object(file_metadata, &mut sink),
594 }
595 }
596
597 fn write_offset_index(
601 &self,
602 offset_index: &OffsetIndexMetaData,
603 column_chunk: &ColumnChunkMetaData,
604 row_group_idx: usize,
605 column_idx: usize,
606 sink: impl Write,
607 ) -> Result<()> {
608 match &self.file_encryptor {
609 Some(file_encryptor) => Self::write_thrift_object_with_encryption(
610 offset_index,
611 sink,
612 file_encryptor,
613 column_chunk,
614 ModuleType::OffsetIndex,
615 row_group_idx,
616 column_idx,
617 ),
618 None => Self::write_thrift_object(offset_index, sink),
619 }
620 }
621
622 fn write_column_index(
628 &self,
629 column_index: &ColumnIndexMetaData,
630 column_chunk: &ColumnChunkMetaData,
631 row_group_idx: usize,
632 column_idx: usize,
633 sink: impl Write,
634 ) -> Result<bool> {
635 match &self.file_encryptor {
636 Some(file_encryptor) => Self::write_thrift_object_with_encryption(
637 column_index,
638 sink,
639 file_encryptor,
640 column_chunk,
641 ModuleType::ColumnIndex,
642 row_group_idx,
643 column_idx,
644 )?,
645 None => Self::write_thrift_object(column_index, sink)?,
646 }
647 Ok(true)
648 }
649
650 fn apply_row_group_encryption(
654 &self,
655 row_groups: Vec<RowGroupMetaData>,
656 ) -> Result<(Vec<RowGroupMetaData>, Option<Vec<RowGroupMetaData>>)> {
657 match &self.file_encryptor {
658 Some(file_encryptor) => {
659 let unencrypted_row_groups = row_groups.clone();
660 let encrypted_row_groups = Self::encrypt_row_groups(row_groups, file_encryptor)?;
661 Ok((encrypted_row_groups, Some(unencrypted_row_groups)))
662 }
663 None => Ok((row_groups, None)),
664 }
665 }
666
667 fn get_file_magic(&self) -> &[u8; 4] {
669 get_file_magic(
670 self.file_encryptor
671 .as_ref()
672 .map(|encryptor| encryptor.properties()),
673 )
674 }
675
676 fn write_thrift_object_with_encryption(
677 object: &impl WriteThrift,
678 mut sink: impl Write,
679 file_encryptor: &FileEncryptor,
680 column_metadata: &ColumnChunkMetaData,
681 module_type: ModuleType,
682 row_group_index: usize,
683 column_index: usize,
684 ) -> Result<()> {
685 let column_path_vec = column_metadata.column_path().as_ref();
686
687 let joined_column_path;
688 let column_path = if column_path_vec.len() == 1 {
689 &column_path_vec[0]
690 } else {
691 joined_column_path = column_path_vec.join(".");
692 &joined_column_path
693 };
694
695 if file_encryptor.is_column_encrypted(column_path) {
696 use crate::encryption::encrypt::encrypt_thrift_object;
697
698 let aad = create_module_aad(
699 file_encryptor.file_aad(),
700 module_type,
701 row_group_index,
702 column_index,
703 None,
704 )?;
705 let mut encryptor = file_encryptor.get_column_encryptor(column_path)?;
706 encrypt_thrift_object(object, &mut encryptor, &mut sink, &aad)
707 } else {
708 Self::write_thrift_object(object, sink)
709 }
710 }
711
712 fn get_plaintext_footer_crypto_metadata(
713 &self,
714 ) -> (Option<EncryptionAlgorithm>, Option<Vec<u8>>) {
715 if let Some(file_encryptor) = self.file_encryptor.as_ref() {
717 let encryption_properties = file_encryptor.properties();
718 if !encryption_properties.encrypt_footer() {
719 return (
720 Some(Self::encryption_algorithm_from_encryptor(file_encryptor)),
721 encryption_properties.footer_key_metadata().cloned(),
722 );
723 }
724 }
725 (None, None)
726 }
727
728 fn encryption_algorithm_from_encryptor(file_encryptor: &FileEncryptor) -> EncryptionAlgorithm {
729 let supply_aad_prefix = file_encryptor
730 .properties()
731 .aad_prefix()
732 .map(|_| !file_encryptor.properties().store_aad_prefix());
733 let aad_prefix = if file_encryptor.properties().store_aad_prefix() {
734 file_encryptor.properties().aad_prefix()
735 } else {
736 None
737 };
738 EncryptionAlgorithm::AES_GCM_V1(AesGcmV1 {
739 aad_prefix: aad_prefix.cloned(),
740 aad_file_unique: Some(file_encryptor.aad_file_unique().clone()),
741 supply_aad_prefix,
742 })
743 }
744
745 fn file_crypto_metadata(file_encryptor: &'_ FileEncryptor) -> Result<FileCryptoMetaData<'_>> {
746 let properties = file_encryptor.properties();
747 Ok(FileCryptoMetaData {
748 encryption_algorithm: Self::encryption_algorithm_from_encryptor(file_encryptor),
749 key_metadata: properties.footer_key_metadata().map(|v| v.as_slice()),
750 })
751 }
752
753 fn encrypt_row_groups(
754 row_groups: Vec<RowGroupMetaData>,
755 file_encryptor: &Arc<FileEncryptor>,
756 ) -> Result<Vec<RowGroupMetaData>> {
757 row_groups
758 .into_iter()
759 .enumerate()
760 .map(|(rg_idx, mut rg)| {
761 let cols: Result<Vec<ColumnChunkMetaData>> = rg
762 .columns
763 .into_iter()
764 .enumerate()
765 .map(|(col_idx, c)| {
766 Self::encrypt_column_chunk(c, file_encryptor, rg_idx, col_idx)
767 })
768 .collect();
769 rg.columns = cols?;
770 Ok(rg)
771 })
772 .collect()
773 }
774
775 fn encrypt_column_chunk(
777 mut column_chunk: ColumnChunkMetaData,
778 file_encryptor: &Arc<FileEncryptor>,
779 row_group_index: usize,
780 column_index: usize,
781 ) -> Result<ColumnChunkMetaData> {
782 let encryptor = match column_chunk.column_crypto_metadata.as_deref() {
785 None => None,
786 Some(ColumnCryptoMetaData::ENCRYPTION_WITH_FOOTER_KEY) => {
787 let is_footer_encrypted = file_encryptor.properties().encrypt_footer();
788
789 if !is_footer_encrypted {
794 Some(file_encryptor.get_footer_encryptor()?)
795 } else {
796 None
797 }
798 }
799 Some(ColumnCryptoMetaData::ENCRYPTION_WITH_COLUMN_KEY(col_key)) => {
800 let column_path = col_key.path_in_schema.join(".");
801 Some(file_encryptor.get_column_encryptor(&column_path)?)
802 }
803 };
804
805 if let Some(mut encryptor) = encryptor {
806 use crate::file::metadata::thrift::serialize_column_meta_data;
807
808 let aad = create_module_aad(
809 file_encryptor.file_aad(),
810 ModuleType::ColumnMetaData,
811 row_group_index,
812 column_index,
813 None,
814 )?;
815 let mut buffer: Vec<u8> = vec![];
817 {
818 let mut prot = ThriftCompactOutputProtocol::new(&mut buffer);
819 serialize_column_meta_data(&column_chunk, &mut prot)?;
820 }
821 let ciphertext = encryptor.encrypt(&buffer, &aad)?;
822 column_chunk.encrypted_column_metadata = Some(ciphertext);
823 column_chunk.plaintext_footer_mode = !file_encryptor.properties().encrypt_footer();
826 }
827
828 Ok(column_chunk)
829 }
830}