laskoviymishka commented on code in PR #3242:
URL: https://github.com/apache/iceberg-rust/pull/3242#discussion_r4102918665
##########
crates/iceberg/src/arrow/reader/pipeline.rs:
##########
@@ -723,6 +635,68 @@ 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.
Review Comment:
The condensed comment keeps "matching Java's ReadConf constructor" but drops
the spec Column Projection quote, the spec URL, and the three Java method names
(`applyNameMapping` / `addFallbackIds` / `pruneColumnsFallback`). That block
was the one spot anchoring branches 2/3 to Java for future cross-client parity
audits — the `hasIds()` note above survived, so only these two lost their
reference. I'd fold the URL and method names back in; condensed is fine.
##########
crates/iceberg/src/arrow/reader/pipeline.rs:
##########
@@ -723,6 +635,68 @@ 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.
+ // The fast path (embedded IDs, no INT96 coercion, no row number)
returns early
+ // without materializing an owned schema.
+ let arrow_schema = if missing_field_ids {
+ let schema = 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())
+ };
+ // 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.
+ coerce_int96_timestamps(&schema, task.schema()).unwrap_or(schema)
+ } else if let Some(coerced) =
+ coerce_int96_timestamps(arrow_metadata.schema(), task.schema())
+ {
+ coerced
+ } else if install_row_number {
+ Arc::clone(arrow_metadata.schema())
+ } else {
+ return Ok(arrow_metadata);
+ };
+
+ let mut options =
ArrowReaderOptions::new().with_schema(Arc::clone(&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,
+ format!(
+ "Failed to create ArrowReaderMetadata with the configured
reader options \
+ (missing_field_ids: {}, install_row_number: {}, schema:
{})",
+ missing_field_ids, install_row_number, arrow_schema,
Review Comment:
These three are plain locals — inline them to match the rest of the file and
the `{coerced_schema}` this replaced. Positional args trip
`clippy::uninlined_format_args`, which bites if CI runs `-D warnings`; the new
test already gets this right with `Row {i}`.
```rust
format!(
"Failed to create ArrowReaderMetadata with the configured reader options
(missing_field_ids: {missing_field_ids}, install_row_number:
{install_row_number}, schema: {arrow_schema})"
)
```
While we're in this message: `arrow_schema` here is the pre-virtual-column
schema, so a failure in the `with_virtual_columns` merge (e.g. a field-id
collision on the synthetic `_pos`) would dump a schema that looks fine. Since
`install_row_number` is already in scope, noting whether that step ran would
close the last corner of the triage detail from round 1.
##########
crates/iceberg/src/arrow/reader/pipeline.rs:
##########
@@ -2847,6 +2821,53 @@ mod tests {
assert_int96_read_matches(&file_path, schema, vec![1, 2],
&expected_micros).await;
}
+ #[tokio::test]
+ async fn test_read_int96_timestamps_with_fallback_ids_and_pos() {
+ use arrow_array::TimestampMicrosecondArray;
+
+ // Regression test for the combined path this refactor introduced: a
field-id-less
+ // file (positional fallback IDs) with an INT96 column and a `_pos`
projection.
+ // All three transforms -- field-ID assignment, INT96 coercion, and
the row-number
+ // virtual column -- apply in the single ArrowReaderMetadata rebuild.
+ let schema = Arc::new(
+ Schema::builder()
+ .with_schema_id(1)
+ .with_fields(vec![
+ NestedField::optional(1, "ts",
Type::Primitive(PrimitiveType::Timestamp))
+ .into(),
+ NestedField::required(2, "id",
Type::Primitive(PrimitiveType::Int)).into(),
+ ])
+ .build()
+ .unwrap(),
+ );
+
+ let tmp_dir = TempDir::new().unwrap();
+ let table_location = tmp_dir.path().to_str().unwrap().to_string();
+ let (file_path, expected_micros) =
+ write_int96_parquet_file(&table_location,
"no_ids_with_pos.parquet", false);
Review Comment:
This lands the combined path from round 1 — thanks. One gap left: `false`
here drives the fallback branch, which is the structurally safe arm
(unconditional ID assign → INT96 → row-number). The arm the collapse most
restructured is the embedded-ID `else if` chain up at line 664, where INT96
coercion sits textually ahead of the row-number check — and nothing exercises
it with both INT96 and a `_pos` projection. A sibling test with
`write_int96_parquet_file(..., true)` would pin it, so a future reorder of
those arms can't silently drop INT96 coercion for embedded-ID files with a row
number. Follow-up, not a blocker — the core combination is covered now.
--
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]