(
&self,
reader: Box<dyn FileRead>,
file_size: u64,
read_fields: &[DataField],
predicates: Option<&FilePredicates>,
batch_size: Option<usize>,
r
| 136 | #[async_trait] |
| 137 | impl FormatFileReader for ParquetFormatReader { |
| 138 | async fn read_batch_stream( |
| 139 | &self, |
| 140 | reader: Box<dyn FileRead>, |
| 141 | file_size: u64, |
| 142 | read_fields: &[DataField], |
| 143 | predicates: Option<&FilePredicates>, |
| 144 | batch_size: Option<usize>, |
| 145 | row_selection: Option<Vec<RowRange>>, |
| 146 | ) -> crate::Result<ArrowRecordBatchStream> { |
| 147 | let arrow_file_reader = ArrowFileReader::new(file_size, reader); |
| 148 | |
| 149 | let mut batch_stream_builder = |
| 150 | ParquetRecordBatchStreamBuilder::new(arrow_file_reader).await?; |
| 151 | |
| 152 | let parquet_schema = batch_stream_builder.parquet_schema().clone(); |
| 153 | let root_schema = parquet_schema.root_schema(); |
| 154 | let root_indices: Vec<usize> = read_fields |
| 155 | .iter() |
| 156 | .filter_map(|f| { |
| 157 | root_schema |
| 158 | .get_fields() |
| 159 | .iter() |
| 160 | .position(|pf| pf.name() == f.name()) |
| 161 | }) |
| 162 | .collect(); |
| 163 | |
| 164 | let mask = ProjectionMask::roots(&parquet_schema, root_indices); |
| 165 | batch_stream_builder = batch_stream_builder.with_projection(mask); |
| 166 | |
| 167 | let empty_predicates = Vec::new(); |
| 168 | let (preds, file_fields): (&[Predicate], &[DataField]) = match predicates { |
| 169 | Some(fp) => (&fp.predicates, &fp.file_fields), |
| 170 | None => (&empty_predicates, &[]), |
| 171 | }; |
| 172 | |
| 173 | let parquet_row_filter = build_parquet_row_filter(&parquet_schema, preds, file_fields)?; |
| 174 | if let Some(f) = parquet_row_filter { |
| 175 | batch_stream_builder = batch_stream_builder.with_row_filter(f); |
| 176 | } |
| 177 | |
| 178 | let predicate_row_selection = build_predicate_row_selection( |
| 179 | batch_stream_builder.metadata().row_groups(), |
| 180 | preds, |
| 181 | file_fields, |
| 182 | )?; |
| 183 | let mut combined_selection = predicate_row_selection; |
| 184 | |
| 185 | if let Some(ref ranges) = row_selection { |
| 186 | let range_selection = |
| 187 | build_row_ranges_selection(batch_stream_builder.metadata().row_groups(), ranges); |
| 188 | combined_selection = |
| 189 | intersect_optional_row_selections(combined_selection, Some(range_selection)); |
| 190 | } |
| 191 | if let Some(sel) = combined_selection { |
| 192 | batch_stream_builder = batch_stream_builder.with_row_selection(sel); |
| 193 | } |
| 194 | if let Some(size) = batch_size { |
| 195 | batch_stream_builder = batch_stream_builder.with_batch_size(size); |
nothing calls this directly
no test coverage detected