1use crate::file::metadata::thrift::FileMeta;
19use crate::file::metadata::{
20 ColumnChunkMetaData, ParquetColumnIndex, ParquetOffsetIndex, RowGroupMetaData,
21};
22use crate::schema::types::{SchemaDescPtr, SchemaDescriptor};
23use crate::{
24 basic::ColumnOrder,
25 file::metadata::{FileMetaData, ParquetMetaDataBuilder},
26};
27#[cfg(feature = "encryption")]
28use crate::{
29 encryption::{
30 encrypt::{FileEncryptor, encrypt_thrift_object, write_signed_plaintext_thrift_object},
31 modules::{ModuleType, create_footer_aad, create_module_aad},
32 },
33 file::column_crypto_metadata::ColumnCryptoMetaData,
34 file::metadata::thrift::encryption::{AesGcmV1, EncryptionAlgorithm, FileCryptoMetaData},
35};
36use crate::{errors::Result, file::page_index::column_index::ColumnIndexMetaData};
37
38use crate::{
39 file::writer::{TrackedWrite, get_file_magic},
40 parquet_thrift::WriteThrift,
41};
42use crate::{
43 file::{
44 metadata::{KeyValue, ParquetMetaData},
45 page_index::offset_index::OffsetIndexMetaData,
46 },
47 parquet_thrift::ThriftCompactOutputProtocol,
48};
49use std::io::Write;
50use std::sync::Arc;
51
52pub(crate) struct ThriftMetadataWriter<'a, W: Write> {
56 buf: &'a mut TrackedWrite<W>,
57 schema_descr: &'a SchemaDescPtr,
58 row_groups: Vec<RowGroupMetaData>,
59 column_indexes: Option<Vec<Vec<Option<ColumnIndexMetaData>>>>,
60 offset_indexes: Option<Vec<Vec<Option<OffsetIndexMetaData>>>>,
61 key_value_metadata: Option<Vec<KeyValue>>,
62 created_by: Option<String>,
63 object_writer: MetadataObjectWriter,
64 writer_version: i32,
65 write_path_in_schema: bool,
66}
67
68impl<'a, W: Write> ThriftMetadataWriter<'a, W> {
69 fn write_offset_indexes(
75 &mut self,
76 offset_indexes: &[Vec<Option<OffsetIndexMetaData>>],
77 ) -> Result<()> {
78 for (row_group_idx, row_group) in self.row_groups.iter_mut().enumerate() {
82 for (column_idx, column_metadata) in row_group.columns.iter_mut().enumerate() {
83 if let Some(offset_index) = &offset_indexes[row_group_idx][column_idx] {
84 let start_offset = self.buf.bytes_written();
85 self.object_writer.write_offset_index(
86 offset_index,
87 column_metadata,
88 row_group_idx,
89 column_idx,
90 &mut self.buf,
91 )?;
92 let end_offset = self.buf.bytes_written();
93 column_metadata.offset_index_offset = Some(start_offset as i64);
95 column_metadata.offset_index_length = Some((end_offset - start_offset) as i32);
96 }
97 }
98 }
99 Ok(())
100 }
101
102 fn write_column_indexes(
108 &mut self,
109 column_indexes: &[Vec<Option<ColumnIndexMetaData>>],
110 ) -> Result<()> {
111 for (row_group_idx, row_group) in self.row_groups.iter_mut().enumerate() {
115 for (column_idx, column_metadata) in row_group.columns.iter_mut().enumerate() {
116 if let Some(column_index) = &column_indexes[row_group_idx][column_idx] {
117 let start_offset = self.buf.bytes_written();
118 if self.object_writer.write_column_index(
120 column_index,
121 column_metadata,
122 row_group_idx,
123 column_idx,
124 &mut self.buf,
125 )? {
126 let end_offset = self.buf.bytes_written();
127 column_metadata.column_index_offset = Some(start_offset as i64);
129 column_metadata.column_index_length =
130 Some((end_offset - start_offset) as i32);
131 }
132 }
133 }
134 }
135 Ok(())
136 }
137
138 fn finalize_column_indexes(&mut self) -> Result<Option<ParquetColumnIndex>> {
140 let column_indexes = std::mem::take(&mut self.column_indexes);
141
142 if let Some(column_indexes) = column_indexes.as_ref() {
144 self.write_column_indexes(column_indexes)?;
145 }
146
147 let all_none = column_indexes
149 .as_ref()
150 .is_some_and(|ci| ci.iter().all(|cii| cii.iter().all(|idx| idx.is_none())));
151
152 let column_indexes: Option<ParquetColumnIndex> = if all_none {
155 None
156 } else {
157 column_indexes.map(|ovvi| {
158 ovvi.into_iter()
159 .map(|vi| {
160 vi.into_iter()
161 .map(|ci| ci.unwrap_or(ColumnIndexMetaData::NONE))
162 .collect()
163 })
164 .collect()
165 })
166 };
167
168 Ok(column_indexes)
169 }
170
171 fn finalize_offset_indexes(&mut self) -> Result<Option<ParquetOffsetIndex>> {
173 let offset_indexes = std::mem::take(&mut self.offset_indexes);
174
175 if let Some(offset_indexes) = offset_indexes.as_ref() {
177 self.write_offset_indexes(offset_indexes)?;
178 }
179
180 let all_none = offset_indexes
182 .as_ref()
183 .is_some_and(|oi| oi.iter().all(|oii| oii.iter().all(|idx| idx.is_none())));
184
185 let offset_indexes: Option<ParquetOffsetIndex> = if all_none {
186 None
187 } else {
188 offset_indexes.map(|ovvi| {
190 ovvi.into_iter()
191 .map(|vi| vi.into_iter().map(|oi| oi.unwrap()).collect())
192 .collect()
193 })
194 };
195
196 Ok(offset_indexes)
197 }
198
199 pub fn finish(mut self) -> Result<ParquetMetaData> {
201 let num_rows = self.row_groups.iter().map(|x| x.num_rows).sum();
202
203 let column_indexes = self.finalize_column_indexes()?;
205 let offset_indexes = self.finalize_offset_indexes()?;
206
207 let column_orders = self
209 .schema_descr
210 .columns()
211 .iter()
212 .map(|col| {
213 ColumnOrder::column_order_for_type(
214 col.logical_type_ref(),
215 col.converted_type(),
216 col.physical_type(),
217 )
218 })
219 .collect();
220
221 let column_orders = Some(column_orders);
225
226 let (row_groups, unencrypted_row_groups) = self
227 .object_writer
228 .apply_row_group_encryption(self.row_groups)?;
229
230 #[cfg(feature = "encryption")]
231 let (encryption_algorithm, footer_signing_key_metadata) =
232 self.object_writer.get_plaintext_footer_crypto_metadata();
233 #[cfg(feature = "encryption")]
234 let file_metadata = FileMetaData::new(
235 self.writer_version,
236 num_rows,
237 self.created_by,
238 self.key_value_metadata,
239 self.schema_descr.clone(),
240 column_orders,
241 )
242 .with_encryption_algorithm(encryption_algorithm)
243 .with_footer_signing_key_metadata(footer_signing_key_metadata);
244
245 #[cfg(not(feature = "encryption"))]
246 let file_metadata = FileMetaData::new(
247 self.writer_version,
248 num_rows,
249 self.created_by,
250 self.key_value_metadata,
251 self.schema_descr.clone(),
252 column_orders,
253 );
254
255 let file_meta = FileMeta {
256 file_metadata: &file_metadata,
257 row_groups: &row_groups,
258 write_path_in_schema: self.write_path_in_schema,
259 };
260
261 let start_pos = self.buf.bytes_written();
263 self.object_writer
264 .write_file_metadata(&file_meta, &mut self.buf)?;
265 let end_pos = self.buf.bytes_written();
266
267 let metadata_len = (end_pos - start_pos) as u32;
269
270 self.buf.write_all(&metadata_len.to_le_bytes())?;
271 self.buf.write_all(self.object_writer.get_file_magic())?;
272
273 let builder = ParquetMetaDataBuilder::new(file_metadata)
278 .set_column_index(column_indexes)
279 .set_offset_index(offset_indexes);
280
281 Ok(match unencrypted_row_groups {
282 Some(rg) => builder.set_row_groups(rg).build(),
283 None => builder.set_row_groups(row_groups).build(),
284 })
285 }
286
287 pub fn new(
288 buf: &'a mut TrackedWrite<W>,
289 schema_descr: &'a SchemaDescPtr,
290 row_groups: Vec<RowGroupMetaData>,
291 created_by: Option<String>,
292 writer_version: i32,
293 write_path_in_schema: bool,
294 ) -> Self {
295 Self {
296 buf,
297 schema_descr,
298 row_groups,
299 column_indexes: None,
300 offset_indexes: None,
301 key_value_metadata: None,
302 created_by,
303 object_writer: Default::default(),
304 writer_version,
305 write_path_in_schema,
306 }
307 }
308
309 pub fn with_column_indexes(
310 mut self,
311 column_indexes: Vec<Vec<Option<ColumnIndexMetaData>>>,
312 ) -> Self {
313 self.column_indexes = Some(column_indexes);
314 self
315 }
316
317 pub fn with_offset_indexes(
318 mut self,
319 offset_indexes: Vec<Vec<Option<OffsetIndexMetaData>>>,
320 ) -> Self {
321 self.offset_indexes = Some(offset_indexes);
322 self
323 }
324
325 pub fn with_key_value_metadata(mut self, key_value_metadata: Vec<KeyValue>) -> Self {
326 self.key_value_metadata = Some(key_value_metadata);
327 self
328 }
329
330 #[cfg(feature = "encryption")]
331 pub fn with_file_encryptor(mut self, file_encryptor: Option<Arc<FileEncryptor>>) -> Self {
332 self.object_writer = self.object_writer.with_file_encryptor(file_encryptor);
333 self
334 }
335}
336
337pub struct ParquetMetaDataWriter<'a, W: Write> {
415 buf: TrackedWrite<W>,
416 metadata: &'a ParquetMetaData,
417 write_path_in_schema: bool,
418}
419
420impl<'a, W: Write> ParquetMetaDataWriter<'a, W> {
421 pub fn new(buf: W, metadata: &'a ParquetMetaData) -> Self {
429 Self::new_with_tracked(TrackedWrite::new(buf), metadata)
430 }
431
432 pub fn new_with_tracked(buf: TrackedWrite<W>, metadata: &'a ParquetMetaData) -> Self {
439 Self {
440 buf,
441 metadata,
442 write_path_in_schema: true,
443 }
444 }
445
446 pub fn with_write_path_in_schema(self, val: bool) -> Self {
449 Self {
450 write_path_in_schema: val,
451 ..self
452 }
453 }
454
455 pub fn finish(mut self) -> Result<()> {
457 let file_metadata = self.metadata.file_metadata();
458
459 let schema = Arc::new(file_metadata.schema().clone());
460 let schema_descr = Arc::new(SchemaDescriptor::new(schema.clone()));
461 let created_by = file_metadata.created_by().map(str::to_string);
462
463 let row_groups = self.metadata.row_groups.clone();
464
465 let key_value_metadata = file_metadata.key_value_metadata().cloned();
466
467 let column_indexes = self.convert_column_indexes();
468 let offset_indexes = self.convert_offset_index();
469
470 let mut encoder = ThriftMetadataWriter::new(
471 &mut self.buf,
472 &schema_descr,
473 row_groups,
474 created_by,
475 file_metadata.version(),
476 self.write_path_in_schema,
477 );
478
479 if let Some(column_indexes) = column_indexes {
480 encoder = encoder.with_column_indexes(column_indexes);
481 }
482
483 if let Some(offset_indexes) = offset_indexes {
484 encoder = encoder.with_offset_indexes(offset_indexes);
485 }
486
487 if let Some(key_value_metadata) = key_value_metadata {
488 encoder = encoder.with_key_value_metadata(key_value_metadata);
489 }
490 encoder.finish()?;
491
492 Ok(())
493 }
494
495 fn convert_column_indexes(&self) -> Option<Vec<Vec<Option<ColumnIndexMetaData>>>> {
496 self.metadata
499 .column_index()
500 .map(|row_group_column_indexes| {
501 (0..self.metadata.row_groups().len())
502 .map(|rg_idx| {
503 let column_indexes = &row_group_column_indexes[rg_idx];
504 column_indexes
505 .iter()
506 .map(|column_index| Some(column_index.clone()))
507 .collect()
508 })
509 .collect()
510 })
511 }
512
513 fn convert_offset_index(&self) -> Option<Vec<Vec<Option<OffsetIndexMetaData>>>> {
514 self.metadata
515 .offset_index()
516 .map(|row_group_offset_indexes| {
517 (0..self.metadata.row_groups().len())
518 .map(|rg_idx| {
519 let offset_indexes = &row_group_offset_indexes[rg_idx];
520 offset_indexes
521 .iter()
522 .map(|offset_index| Some(offset_index.clone()))
523 .collect()
524 })
525 .collect()
526 })
527 }
528}
529
530#[derive(Debug, Default)]
531struct MetadataObjectWriter {
532 #[cfg(feature = "encryption")]
533 file_encryptor: Option<Arc<FileEncryptor>>,
534}
535
536impl MetadataObjectWriter {
537 #[inline]
538 fn write_thrift_object(object: &impl WriteThrift, sink: impl Write) -> Result<()> {
539 let mut protocol = ThriftCompactOutputProtocol::new(sink);
540 object.write_thrift(&mut protocol)?;
541 Ok(())
542 }
543}
544
545#[cfg(not(feature = "encryption"))]
547impl MetadataObjectWriter {
548 fn write_file_metadata(&self, file_metadata: &FileMeta, sink: impl Write) -> Result<()> {
552 Self::write_thrift_object(file_metadata, sink)
553 }
554
555 fn write_offset_index(
559 &self,
560 offset_index: &OffsetIndexMetaData,
561 _column_chunk: &ColumnChunkMetaData,
562 _row_group_idx: usize,
563 _column_idx: usize,
564 sink: impl Write,
565 ) -> Result<()> {
566 Self::write_thrift_object(offset_index, sink)
567 }
568
569 fn write_column_index(
576 &self,
577 column_index: &ColumnIndexMetaData,
578 _column_chunk: &ColumnChunkMetaData,
579 _row_group_idx: usize,
580 _column_idx: usize,
581 sink: impl Write,
582 ) -> Result<bool> {
583 match column_index {
584 ColumnIndexMetaData::NONE => Ok(false),
586 _ => {
587 Self::write_thrift_object(column_index, sink)?;
588 Ok(true)
589 }
590 }
591 }
592
593 fn apply_row_group_encryption(
595 &self,
596 row_groups: Vec<RowGroupMetaData>,
597 ) -> Result<(Vec<RowGroupMetaData>, Option<Vec<RowGroupMetaData>>)> {
598 Ok((row_groups, None))
599 }
600
601 pub fn get_file_magic(&self) -> &[u8; 4] {
603 get_file_magic()
604 }
605}
606
607#[cfg(feature = "encryption")]
609impl MetadataObjectWriter {
610 fn with_file_encryptor(mut self, encryptor: Option<Arc<FileEncryptor>>) -> Self {
612 self.file_encryptor = encryptor;
613 self
614 }
615
616 fn write_file_metadata(&self, file_metadata: &FileMeta, mut sink: impl Write) -> Result<()> {
620 match self.file_encryptor.as_ref() {
621 Some(file_encryptor) if file_encryptor.properties().encrypt_footer() => {
622 let crypto_metadata = Self::file_crypto_metadata(file_encryptor)?;
624 let mut protocol = ThriftCompactOutputProtocol::new(&mut sink);
625 crypto_metadata.write_thrift(&mut protocol)?;
626
627 let aad = create_footer_aad(file_encryptor.file_aad())?;
629 let mut encryptor = file_encryptor.get_footer_encryptor()?;
630 encrypt_thrift_object(file_metadata, &mut encryptor, &mut sink, &aad)
631 }
632 Some(file_encryptor) if file_metadata.file_metadata.encryption_algorithm.is_some() => {
633 let aad = create_footer_aad(file_encryptor.file_aad())?;
634 let mut encryptor = file_encryptor.get_footer_encryptor()?;
635 write_signed_plaintext_thrift_object(file_metadata, &mut encryptor, &mut sink, &aad)
636 }
637 _ => Self::write_thrift_object(file_metadata, &mut sink),
638 }
639 }
640
641 fn write_offset_index(
645 &self,
646 offset_index: &OffsetIndexMetaData,
647 column_chunk: &ColumnChunkMetaData,
648 row_group_idx: usize,
649 column_idx: usize,
650 sink: impl Write,
651 ) -> Result<()> {
652 match &self.file_encryptor {
653 Some(file_encryptor) => Self::write_thrift_object_with_encryption(
654 offset_index,
655 sink,
656 file_encryptor,
657 column_chunk,
658 ModuleType::OffsetIndex,
659 row_group_idx,
660 column_idx,
661 ),
662 None => Self::write_thrift_object(offset_index, sink),
663 }
664 }
665
666 fn write_column_index(
673 &self,
674 column_index: &ColumnIndexMetaData,
675 column_chunk: &ColumnChunkMetaData,
676 row_group_idx: usize,
677 column_idx: usize,
678 sink: impl Write,
679 ) -> Result<bool> {
680 match column_index {
681 ColumnIndexMetaData::NONE => Ok(false),
683 _ => {
684 match &self.file_encryptor {
685 Some(file_encryptor) => Self::write_thrift_object_with_encryption(
686 column_index,
687 sink,
688 file_encryptor,
689 column_chunk,
690 ModuleType::ColumnIndex,
691 row_group_idx,
692 column_idx,
693 )?,
694 None => Self::write_thrift_object(column_index, sink)?,
695 }
696 Ok(true)
697 }
698 }
699 }
700
701 fn apply_row_group_encryption(
705 &self,
706 row_groups: Vec<RowGroupMetaData>,
707 ) -> Result<(Vec<RowGroupMetaData>, Option<Vec<RowGroupMetaData>>)> {
708 match &self.file_encryptor {
709 Some(file_encryptor) => {
710 let unencrypted_row_groups = row_groups.clone();
711 let encrypted_row_groups = Self::encrypt_row_groups(row_groups, file_encryptor)?;
712 Ok((encrypted_row_groups, Some(unencrypted_row_groups)))
713 }
714 None => Ok((row_groups, None)),
715 }
716 }
717
718 fn get_file_magic(&self) -> &[u8; 4] {
720 get_file_magic(
721 self.file_encryptor
722 .as_ref()
723 .map(|encryptor| encryptor.properties()),
724 )
725 }
726
727 fn write_thrift_object_with_encryption(
728 object: &impl WriteThrift,
729 mut sink: impl Write,
730 file_encryptor: &FileEncryptor,
731 column_metadata: &ColumnChunkMetaData,
732 module_type: ModuleType,
733 row_group_index: usize,
734 column_index: usize,
735 ) -> Result<()> {
736 let column_path_vec = column_metadata.column_path().as_ref();
737
738 let joined_column_path;
739 let column_path = if column_path_vec.len() == 1 {
740 &column_path_vec[0]
741 } else {
742 joined_column_path = column_path_vec.join(".");
743 &joined_column_path
744 };
745
746 if file_encryptor.is_column_encrypted(column_path) {
747 use crate::encryption::encrypt::encrypt_thrift_object;
748
749 let aad = create_module_aad(
750 file_encryptor.file_aad(),
751 module_type,
752 row_group_index,
753 column_index,
754 None,
755 )?;
756 let mut encryptor = file_encryptor.get_column_encryptor(column_path)?;
757 encrypt_thrift_object(object, &mut encryptor, &mut sink, &aad)
758 } else {
759 Self::write_thrift_object(object, sink)
760 }
761 }
762
763 fn get_plaintext_footer_crypto_metadata(
764 &self,
765 ) -> (Option<EncryptionAlgorithm>, Option<Vec<u8>>) {
766 if let Some(file_encryptor) = self.file_encryptor.as_ref() {
768 let encryption_properties = file_encryptor.properties();
769 if !encryption_properties.encrypt_footer() {
770 return (
771 Some(Self::encryption_algorithm_from_encryptor(file_encryptor)),
772 encryption_properties.footer_key_metadata().cloned(),
773 );
774 }
775 }
776 (None, None)
777 }
778
779 fn encryption_algorithm_from_encryptor(file_encryptor: &FileEncryptor) -> EncryptionAlgorithm {
780 let supply_aad_prefix = file_encryptor
781 .properties()
782 .aad_prefix()
783 .map(|_| !file_encryptor.properties().store_aad_prefix());
784 let aad_prefix = if file_encryptor.properties().store_aad_prefix() {
785 file_encryptor.properties().aad_prefix()
786 } else {
787 None
788 };
789 EncryptionAlgorithm::AES_GCM_V1(AesGcmV1 {
790 aad_prefix: aad_prefix.cloned(),
791 aad_file_unique: Some(file_encryptor.aad_file_unique().clone()),
792 supply_aad_prefix,
793 })
794 }
795
796 fn file_crypto_metadata(file_encryptor: &'_ FileEncryptor) -> Result<FileCryptoMetaData<'_>> {
797 let properties = file_encryptor.properties();
798 Ok(FileCryptoMetaData {
799 encryption_algorithm: Self::encryption_algorithm_from_encryptor(file_encryptor),
800 key_metadata: properties.footer_key_metadata().map(|v| v.as_slice()),
801 })
802 }
803
804 fn encrypt_row_groups(
805 row_groups: Vec<RowGroupMetaData>,
806 file_encryptor: &Arc<FileEncryptor>,
807 ) -> Result<Vec<RowGroupMetaData>> {
808 row_groups
809 .into_iter()
810 .enumerate()
811 .map(|(rg_idx, mut rg)| {
812 let cols: Result<Vec<ColumnChunkMetaData>> = rg
813 .columns
814 .into_iter()
815 .enumerate()
816 .map(|(col_idx, c)| {
817 Self::encrypt_column_chunk(c, file_encryptor, rg_idx, col_idx)
818 })
819 .collect();
820 rg.columns = cols?;
821 Ok(rg)
822 })
823 .collect()
824 }
825
826 fn encrypt_column_chunk(
828 mut column_chunk: ColumnChunkMetaData,
829 file_encryptor: &Arc<FileEncryptor>,
830 row_group_index: usize,
831 column_index: usize,
832 ) -> Result<ColumnChunkMetaData> {
833 let encryptor = match column_chunk.column_crypto_metadata.as_deref() {
836 None => None,
837 Some(ColumnCryptoMetaData::ENCRYPTION_WITH_FOOTER_KEY) => {
838 let is_footer_encrypted = file_encryptor.properties().encrypt_footer();
839
840 if !is_footer_encrypted {
845 Some(file_encryptor.get_footer_encryptor()?)
846 } else {
847 None
848 }
849 }
850 Some(ColumnCryptoMetaData::ENCRYPTION_WITH_COLUMN_KEY(col_key)) => {
851 let column_path = col_key.path_in_schema.join(".");
852 Some(file_encryptor.get_column_encryptor(&column_path)?)
853 }
854 };
855
856 if let Some(mut encryptor) = encryptor {
857 use crate::file::metadata::thrift::serialize_column_meta_data;
858
859 let aad = create_module_aad(
860 file_encryptor.file_aad(),
861 ModuleType::ColumnMetaData,
862 row_group_index,
863 column_index,
864 None,
865 )?;
866 let mut buffer: Vec<u8> = vec![];
868 {
869 let mut prot = ThriftCompactOutputProtocol::new(&mut buffer);
870 serialize_column_meta_data(&column_chunk, &mut prot)?;
871 }
872 let ciphertext = encryptor.encrypt(&buffer, &aad)?;
873 column_chunk.encrypted_column_metadata = Some(ciphertext);
874 column_chunk.plaintext_footer_mode = !file_encryptor.properties().encrypt_footer();
877 }
878
879 Ok(column_chunk)
880 }
881}