use std::sync::Arc;
use crate::errors::ParquetError;
use crate::file::metadata::page_index::{PageIndexBuilder, PageIndexProvider};
use crate::file::metadata::thrift::parquet_metadata_from_bytes;
use crate::file::metadata::{
ColumnChunkMetaData, PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions,
};
use crate::file::page_index::column_index::ColumnIndexMetaData;
use crate::file::page_index::index_reader::{decode_column_index, decode_offset_index};
use crate::file::page_index::offset_index::OffsetIndexMetaData;
use bytes::Bytes;
pub(crate) use inner::MetadataParser;
#[cfg(feature = "encryption")]
mod inner {
use std::sync::Arc;
use super::*;
use crate::encryption::decrypt::FileDecryptionProperties;
use crate::errors::Result;
#[derive(Debug, Default)]
pub(crate) struct MetadataParser {
file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
metadata_options: Option<Arc<ParquetMetaDataOptions>>,
}
impl MetadataParser {
pub(crate) fn new() -> Self {
MetadataParser::default()
}
pub(crate) fn with_file_decryption_properties(
mut self,
file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
) -> Self {
self.file_decryption_properties = file_decryption_properties;
self
}
pub(crate) fn with_metadata_options(
self,
options: Option<Arc<ParquetMetaDataOptions>>,
) -> Self {
Self {
metadata_options: options,
..self
}
}
pub(crate) fn decode_metadata(
&self,
buf: &[u8],
encrypted_footer: bool,
) -> Result<ParquetMetaData> {
if encrypted_footer || self.file_decryption_properties.is_some() {
crate::file::metadata::thrift::encryption::parquet_metadata_with_encryption(
self.file_decryption_properties.as_ref(),
encrypted_footer,
buf,
self.metadata_options.as_deref(),
)
} else {
decode_metadata(buf, self.metadata_options.as_deref())
}
}
}
pub(super) fn parse_single_column_index(
bytes: &[u8],
metadata: &ParquetMetaData,
column: &ColumnChunkMetaData,
row_group_index: usize,
col_index: usize,
) -> crate::errors::Result<ColumnIndexMetaData> {
use crate::encryption::decrypt::CryptoContext;
match &column.column_crypto_metadata {
Some(crypto_metadata) => {
let file_decryptor = metadata.file_decryptor.as_ref().ok_or_else(|| {
general_err!("Cannot decrypt column index, no file decryptor set")
})?;
let crypto_context = CryptoContext::for_column(
file_decryptor,
crypto_metadata,
row_group_index,
col_index,
)?;
let column_decryptor = crypto_context.metadata_decryptor();
let aad = crypto_context.create_column_index_aad()?;
let plaintext = column_decryptor.decrypt(bytes, &aad)?;
decode_column_index(&plaintext, column.column_type())
}
None => decode_column_index(bytes, column.column_type()),
}
}
pub(super) fn parse_single_offset_index(
bytes: &[u8],
metadata: &ParquetMetaData,
column: &ColumnChunkMetaData,
row_group_index: usize,
col_index: usize,
) -> crate::errors::Result<OffsetIndexMetaData> {
use crate::encryption::decrypt::CryptoContext;
match &column.column_crypto_metadata {
Some(crypto_metadata) => {
let file_decryptor = metadata.file_decryptor.as_ref().ok_or_else(|| {
general_err!("Cannot decrypt offset index, no file decryptor set")
})?;
let crypto_context = CryptoContext::for_column(
file_decryptor,
crypto_metadata,
row_group_index,
col_index,
)?;
let column_decryptor = crypto_context.metadata_decryptor();
let aad = crypto_context.create_offset_index_aad()?;
let plaintext = column_decryptor.decrypt(bytes, &aad)?;
decode_offset_index(&plaintext)
}
None => decode_offset_index(bytes),
}
}
}
#[cfg(not(feature = "encryption"))]
mod inner {
use super::*;
use crate::errors::Result;
use std::sync::Arc;
#[derive(Debug, Default)]
pub(crate) struct MetadataParser {
metadata_options: Option<Arc<ParquetMetaDataOptions>>,
}
impl MetadataParser {
pub(crate) fn new() -> Self {
MetadataParser::default()
}
pub(crate) fn with_metadata_options(
self,
options: Option<Arc<ParquetMetaDataOptions>>,
) -> Self {
Self {
metadata_options: options,
}
}
pub(crate) fn decode_metadata(
&self,
buf: &[u8],
encrypted_footer: bool,
) -> Result<ParquetMetaData> {
if encrypted_footer {
Err(general_err!(
"Parquet file has an encrypted footer but the encryption feature is disabled"
))
} else {
decode_metadata(buf, self.metadata_options.as_deref())
}
}
}
pub(super) fn parse_single_column_index(
bytes: &[u8],
_metadata: &ParquetMetaData,
column: &ColumnChunkMetaData,
_row_group_index: usize,
_col_index: usize,
) -> crate::errors::Result<ColumnIndexMetaData> {
decode_column_index(bytes, column.column_type())
}
pub(super) fn parse_single_offset_index(
bytes: &[u8],
_metadata: &ParquetMetaData,
_column: &ColumnChunkMetaData,
_row_group_index: usize,
_col_index: usize,
) -> crate::errors::Result<OffsetIndexMetaData> {
decode_offset_index(bytes)
}
}
pub(crate) fn decode_metadata(
buf: &[u8],
options: Option<&ParquetMetaDataOptions>,
) -> crate::errors::Result<ParquetMetaData> {
parquet_metadata_from_bytes(buf, options)
}
pub(crate) fn parse_page_index(
metadata: &mut ParquetMetaData,
column_index_policy: PageIndexPolicy,
offset_index_policy: PageIndexPolicy,
bytes: &Bytes,
start_offset: u64,
) -> crate::errors::Result<()> {
if column_index_policy == PageIndexPolicy::Skip && offset_index_policy == PageIndexPolicy::Skip
{
return Ok(());
}
let num_row_groups = metadata.num_row_groups();
let num_columns = metadata.file_metadata().schema_descr().num_columns();
let mut builder = PageIndexBuilder::default();
if column_index_policy != PageIndexPolicy::Skip {
builder.allocate_column_indexes(num_row_groups, num_columns);
parse_column_index(
metadata,
column_index_policy,
&mut builder,
bytes,
start_offset,
)?;
}
if offset_index_policy != PageIndexPolicy::Skip {
builder.allocate_offset_indexes(num_row_groups, num_columns);
parse_offset_index(
metadata,
offset_index_policy,
&mut builder,
bytes,
start_offset,
)?;
}
let page_index = builder.build();
if !page_index.has_column_indexes() && !page_index.has_offset_indexes() {
return Ok(());
}
metadata.set_page_index(Some(Arc::new(page_index)));
Ok(())
}
fn parse_column_index(
metadata: &ParquetMetaData,
column_index_policy: PageIndexPolicy,
page_index_builder: &mut PageIndexBuilder,
bytes: &Bytes,
start_offset: u64,
) -> crate::errors::Result<()> {
if column_index_policy == PageIndexPolicy::Skip {
return Ok(());
}
for rg_idx in 0..metadata.num_row_groups() {
let rg = metadata.row_group(rg_idx);
for col_idx in 0..rg.num_columns() {
let col = rg.column(col_idx);
if let Some(r) = col.column_index_range() {
let r_start = usize::try_from(r.start - start_offset)?;
let r_end = usize::try_from(r.end - start_offset)?;
let idx = inner::parse_single_column_index(
&bytes[r_start..r_end],
metadata,
col,
rg_idx,
col_idx,
)?;
page_index_builder.put_column_index(idx, rg_idx, col_idx);
}
}
}
Ok(())
}
fn parse_offset_index(
metadata: &ParquetMetaData,
offset_index_policy: PageIndexPolicy,
page_index_builder: &mut PageIndexBuilder,
bytes: &Bytes,
start_offset: u64,
) -> crate::errors::Result<()> {
if offset_index_policy == PageIndexPolicy::Skip {
return Ok(());
}
for rg_idx in 0..metadata.num_row_groups() {
let rg = metadata.row_group(rg_idx);
for col_idx in 0..rg.num_columns() {
let col = rg.column(col_idx);
if let Some(r) = col.offset_index_range() {
let r_start = usize::try_from(r.start - start_offset)?;
let r_end = usize::try_from(r.end - start_offset)?;
let idx = inner::parse_single_offset_index(
&bytes[r_start..r_end],
metadata,
col,
rg_idx,
col_idx,
)?;
page_index_builder.put_offset_index(idx, rg_idx, col_idx);
} else if offset_index_policy == PageIndexPolicy::Required {
return Err(general_err!("missing offset index"));
}
}
}
Ok(())
}