This is an automated email from the ASF dual-hosted git repository.
Jefffrey pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new 78b40bf2f6 GH-10382: Fix CSV RecordBatchReader projected schema
(#10611)
78b40bf2f6 is described below
commit 78b40bf2f62121c56b7c358bdf768f769ebc7581
Author: Langning Zhang <[email protected]>
AuthorDate: Wed Aug 12 15:44:36 2026 +0800
GH-10382: Fix CSV RecordBatchReader projected schema (#10611)
# Which issue does this PR close?
- Closes #10382.
# Rationale for this change
CSV readers configured with a projection yielded projected `RecordBatch`
values, but the `RecordBatchReader::schema` implementation on
`BufReader` still reported the unprojected input schema. This made the
reader's advertised schema disagree with the batches it produced.
# What changes are included in this PR?
- Reuse a shared projected-schema helper for both CSV reader schema
APIs.
- Build projected schemas with `Schema::project`.
- Add regression coverage for no projection, an empty projection, a
single-column projection, and reordered columns.
# Are these changes tested?
Yes. The following checks passed locally:
- `cargo test -p arrow-csv --all-features`
- `cargo fmt --all -- --check`
- `cargo clippy -p arrow-csv --all-targets --all-features -- -D
warnings`
- `cargo clippy -p arrow-csv -- -D unused_crate_dependencies`
- `cargo clippy -p arrow-csv --all-features -- -D
unused_crate_dependencies`
- `cargo clippy -p arrow-csv --no-default-features -- -D
unused_crate_dependencies`
# Are there any user-facing changes?
Yes. `RecordBatchReader::schema` now reports the projected schema for
projected CSV readers, matching the schema of the record batches they
produce. There is no public API change.
# AI usage
OpenAI Codex assisted with implementing the shared projected-schema
logic and regression test. I reviewed every generated change and can
debug and maintain it.
---
arrow-csv/src/reader/mod.rs | 101 +++++++++++++++++++++++++++++++++++---------
1 file changed, 81 insertions(+), 20 deletions(-)
diff --git a/arrow-csv/src/reader/mod.rs b/arrow-csv/src/reader/mod.rs
index 9007acc14e..fc06aaffa5 100644
--- a/arrow-csv/src/reader/mod.rs
+++ b/arrow-csv/src/reader/mod.rs
@@ -506,6 +506,8 @@ pub struct BufReader<R> {
reader: R,
/// The decoder
decoder: Decoder,
+ /// Schema of the record batches produced by this reader
+ schema: SchemaRef,
}
impl<R> fmt::Debug for BufReader<R>
@@ -563,14 +565,7 @@ impl<R: Read> Reader<R> {
/// Returns the schema of the reader, useful for getting the schema
without reading
/// record batches
pub fn schema(&self) -> SchemaRef {
- match &self.decoder.projection {
- Some(projection) => {
- let fields = self.decoder.schema.fields();
- let projected = projection.iter().map(|i| fields[*i].clone());
- Arc::new(Schema::new(projected.collect::<Fields>()))
- }
- None => self.decoder.schema.clone(),
- }
+ self.schema.clone()
}
}
@@ -604,7 +599,7 @@ impl<R: BufRead> Iterator for BufReader<R> {
impl<R: BufRead> RecordBatchReader for BufReader<R> {
fn schema(&self) -> SchemaRef {
- self.decoder.schema.clone()
+ self.schema.clone()
}
}
@@ -730,8 +725,7 @@ impl Decoder {
let rows = self.record_decoder.flush()?;
let batch = parse(
&rows,
- self.schema.fields(),
- Some(self.schema.metadata.clone()),
+ &self.schema,
self.projection.as_ref(),
self.line_number,
&self.null_regex,
@@ -784,16 +778,17 @@ fn validate_header(rows: &StringRecords<'_>, fields:
&Fields) -> Result<(), Arro
/// Parses a slice of [`StringRecords`] into a [RecordBatch]
fn parse(
rows: &StringRecords<'_>,
- fields: &Fields,
- metadata: Option<Metadata>,
+ schema: &Schema,
projection: Option<&Vec<usize>>,
line_number: usize,
null_regex: &NullRegex,
) -> Result<RecordBatch, ArrowError> {
+ let fields = schema.fields();
let projection: Vec<usize> = match projection {
Some(v) => v.clone(),
None => fields.iter().enumerate().map(|(i, _)| i).collect(),
};
+ let projected_schema = Arc::new(schema.project(&projection)?);
let arrays: Result<Vec<ArrayRef>, _> = projection
.iter()
@@ -1022,13 +1017,6 @@ fn parse(
})
.collect();
- let projected_fields: Fields = projection.iter().map(|i|
fields[*i].clone()).collect();
-
- let projected_schema = Arc::new(match metadata {
- None => Schema::new(projected_fields),
- Some(metadata) => Schema::new_with_metadata(projected_fields,
metadata),
- });
-
arrays.and_then(|arr| {
RecordBatch::try_new_with_options(
projected_schema,
@@ -1350,9 +1338,15 @@ impl ReaderBuilder {
/// Create a new `BufReader` from a buffered reader
pub fn build_buffered<R: BufRead>(self, reader: R) -> Result<BufReader<R>,
ArrowError> {
+ let schema = match &self.projection {
+ Some(projection) => Arc::new(self.schema.project(projection)?),
+ None => self.schema.clone(),
+ };
+
Ok(BufReader {
reader,
decoder: self.build_decoder(),
+ schema,
})
}
@@ -1700,6 +1694,73 @@ mod tests {
assert_eq!(2, batch.num_columns());
}
+ #[test]
+ fn test_csv_record_batch_reader_schema() {
+ let schema = Arc::new(Schema::new(vec![
+ Field::new("a", DataType::Int32, false),
+ Field::new("b", DataType::Int32, false),
+ ]));
+
+ let cases = [
+ None,
+ Some(vec![]),
+ Some(vec![1]),
+ Some(vec![1, 0]),
+ Some(vec![1, 1]),
+ ];
+ for projection in cases {
+ let builder = ReaderBuilder::new(schema.clone());
+ let builder = match projection {
+ Some(projection) => builder.with_projection(projection),
+ None => builder,
+ };
+ let mut reader = builder.build(Cursor::new(b"1,2\n")).unwrap();
+
+ let reader_schema = RecordBatchReader::schema(&reader);
+ let batch = reader.next().unwrap().unwrap();
+
+ assert_eq!(reader_schema, batch.schema());
+ }
+ }
+
+ #[test]
+ fn test_csv_reader_rejects_invalid_projection() {
+ let schema = Arc::new(Schema::new(vec![
+ Field::new("a", DataType::Int32, false),
+ Field::new("b", DataType::Int32, false),
+ ]));
+
+ let result = ReaderBuilder::new(schema)
+ .with_projection(vec![2])
+ .build(Cursor::new(b"1,2\n"));
+
+ assert!(matches!(
+ result,
+ Err(ArrowError::SchemaError(message))
+ if message == "project index 2 out of bounds, max field 2"
+ ));
+ }
+
+ #[test]
+ fn test_csv_decoder_rejects_invalid_projection() {
+ let schema = Arc::new(Schema::new(vec![
+ Field::new("a", DataType::Int32, false),
+ Field::new("b", DataType::Int32, false),
+ ]));
+ let mut decoder = ReaderBuilder::new(schema)
+ .with_projection(vec![2])
+ .build_decoder();
+
+ decoder.decode(b"1,2\n").unwrap();
+ let result = decoder.flush();
+
+ assert!(matches!(
+ result,
+ Err(ArrowError::SchemaError(message))
+ if message == "project index 2 out of bounds, max field 2"
+ ));
+ }
+
#[test]
fn test_csv_with_dictionary() {
let schema = Arc::new(Schema::new(vec![