kamcheungting-db commented on code in PR #3242:
URL: https://github.com/apache/iceberg-rust/pull/3242#discussion_r4141847770
##########
crates/iceberg/src/arrow/reader/pipeline.rs:
##########
@@ -723,6 +635,65 @@ impl FileScanTaskReader {
Ok(Box::pin(record_batch_stream) as ArrowRecordBatchStream)
}
+ /// Applies all task-specific schema and virtual-column options,
rebuilding the
+ /// Arrow reader metadata at most once.
+ fn configure_arrow_reader_metadata(
+ arrow_metadata: ArrowReaderMetadata,
+ task: &FileScanTask,
+ missing_field_ids: bool,
+ install_row_number: bool,
+ ) -> Result<ArrowReaderMetadata> {
+ // Three-branch schema resolution strategy matching Java's ReadConf
constructor.
+ // When Parquet files lack field IDs, apply a name mapping when
available and use
+ // position-based fallback IDs otherwise. Files with embedded IDs keep
their schema.
+ let mut arrow_schema = if missing_field_ids {
+ if let Some(name_mapping) = task.name_mapping() {
+ apply_name_mapping_to_arrow_schema(
+ Arc::clone(arrow_metadata.schema()),
+ name_mapping,
+ )?
+ } else {
+ add_fallback_field_ids_to_arrow_schema(arrow_metadata.schema())
+ }
+ } else {
+ Arc::clone(arrow_metadata.schema())
+ };
+
+ // Coerce INT96 timestamp columns before building the stream reader to
avoid i64
+ // overflow in arrow-rs. Apply this after assigning any missing field
IDs so the
+ // final schema contains both changes.
+ let mut should_rebuild = missing_field_ids;
+ if let Some(coerced_schema) = coerce_int96_timestamps(&arrow_schema,
task.schema()) {
+ arrow_schema = coerced_schema;
+ should_rebuild = true;
+ }
+
+ if !should_rebuild && !install_row_number {
+ return Ok(arrow_metadata);
+ }
+
+ let mut options = ArrowReaderOptions::new().with_schema(arrow_schema);
+ if install_row_number {
+ let row_number_field = Arc::new(
+ Field::new(RESERVED_COL_NAME_POS, DataType::Int64, false)
+ .with_metadata(HashMap::from([(
+ PARQUET_FIELD_ID_META_KEY.to_string(),
+ RESERVED_FIELD_ID_POS.to_string(),
+ )]))
+ .with_extension_type(RowNumber),
+ );
+ options = options.with_virtual_columns(vec![row_number_field])?;
+ }
+
+ ArrowReaderMetadata::try_new(Arc::clone(arrow_metadata.metadata()),
options).map_err(|e| {
+ Error::new(
+ ErrorKind::Unexpected,
+ "Failed to create ArrowReaderMetadata with the configured
reader options",
Review Comment:
Done
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]