Skip to content

Commit d55f945

Browse files
getChanjecsand838
andauthored
Apply suggestions from code review
Co-authored-by: Connor Sanders <170039284+jecsand838@users.noreply.github.com>
1 parent b2e9eb1 commit d55f945

2 files changed

Lines changed: 264 additions & 12 deletions

File tree

datafusion/datasource-avro/src/mod.rs

Lines changed: 176 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,187 @@ use arrow::datatypes::Schema;
3333
pub use arrow_avro;
3434
use arrow_avro::reader::ReaderBuilder;
3535
pub use file_format::*;
36+
use arrow_avro::schema::SCHEMA_METADATA_KEY;
37+
use datafusion_common::DataFusionError;
3638
use std::io::{BufReader, Read};
39+
use std::sync::Arc;
3740

3841
/// Read Avro schema given a reader
3942
pub fn read_avro_schema_from_reader<R: Read>(
4043
reader: &mut R,
4144
) -> datafusion_common::Result<Schema> {
4245
let avro_reader = ReaderBuilder::new().build(BufReader::new(reader))?;
43-
Ok(avro_reader.schema().as_ref().clone())
46+
let schema_ref = avro_reader.schema();
47+
// Extract the raw Avro JSON schema from the OCF header.
48+
let raw_json = avro_reader
49+
.avro_header()
50+
.get(SCHEMA_METADATA_KEY.as_bytes())
51+
.map(|bytes| {
52+
std::str::from_utf8(bytes).map_err(|e| {
53+
DataFusionError::Execution(format!(
54+
"Invalid UTF-8 in Avro schema metadata ({SCHEMA_METADATA_KEY}): {e}"
55+
))
56+
})
57+
})
58+
.transpose()?
59+
.map(str::to_owned);
60+
drop(avro_reader);
61+
if let Some(raw_json) = raw_json {
62+
let mut schema = Arc::unwrap_or_clone(schema_ref);
63+
// Insert the raw Avro JSON schema using `SCHEMA_METADATA_KEY`.
64+
// This should enable the avro schema metadata to be picked downstream.
65+
schema
66+
.metadata
67+
.insert(SCHEMA_METADATA_KEY.to_string(), raw_json);
68+
Ok(schema)
69+
} else {
70+
// Return error because Avro spec requires the Avro schema metadata to be present in the OCF header.
71+
Err(DataFusionError::Execution(format!(
72+
"Avro schema metadata ({SCHEMA_METADATA_KEY}) is missing from OCF header"
73+
)))
74+
}
75+
}
76+
77+
#[cfg(test)]
78+
mod test {
79+
use super::*;
80+
use arrow::array::{BinaryArray, BooleanArray, Float64Array};
81+
use arrow::datatypes::DataType;
82+
use arrow_avro::reader::ReaderBuilder;
83+
use arrow_avro::schema::{AvroSchema, SCHEMA_METADATA_KEY};
84+
use datafusion_common::test_util::arrow_test_data;
85+
use datafusion_common::{DataFusionError, Result as DFResult};
86+
use serde_json::Value;
87+
use std::collections::HashMap;
88+
use std::fs::File;
89+
use std::io::BufReader;
90+
91+
fn avro_test_file(name: &str) -> String {
92+
format!("{}/avro/{name}", arrow_test_data())
93+
}
94+
95+
#[test]
96+
fn read_avro_schema_includes_avro_json_metadata() -> DFResult<()> {
97+
let path = avro_test_file("alltypes_plain.avro");
98+
let mut file = File::open(&path)?;
99+
let schema = read_avro_schema_from_reader(&mut file)?;
100+
let meta_json = schema
101+
.metadata()
102+
.get(SCHEMA_METADATA_KEY)
103+
.expect("schema metadata missing avro.schema entry");
104+
assert!(
105+
!meta_json.is_empty(),
106+
"avro.schema metadata should not be empty"
107+
);
108+
let mut raw = File::open(&path)?;
109+
let avro_reader = ReaderBuilder::new().build(BufReader::new(&mut raw))?;
110+
let header_json = avro_reader
111+
.avro_header()
112+
.get(SCHEMA_METADATA_KEY.as_bytes())
113+
.and_then(|bytes| std::str::from_utf8(bytes).ok())
114+
.expect("missing avro.schema metadata in OCF header");
115+
assert_eq!(
116+
meta_json, header_json,
117+
"schema metadata avro.schema should match OCF header"
118+
);
119+
Ok(())
120+
}
121+
122+
#[test]
123+
fn read_and_project_using_schema_metadata() -> DFResult<()> {
124+
let path = avro_test_file("alltypes_dictionary.avro");
125+
let mut file = File::open(&path)?;
126+
let file_schema = read_avro_schema_from_reader(&mut file)?;
127+
let projected_field_names = vec!["string_col", "double_col", "bool_col"];
128+
let avro_json = file_schema
129+
.metadata()
130+
.get(SCHEMA_METADATA_KEY)
131+
.expect("schema metadata missing avro.schema entry");
132+
let projected_avro_schema =
133+
build_projected_reader_schema(avro_json, &projected_field_names)?;
134+
let mut reader = ReaderBuilder::new()
135+
.with_reader_schema(projected_avro_schema)
136+
.with_batch_size(64)
137+
.build(BufReader::new(File::open(&path)?))?;
138+
let batch = reader.next().expect("no batch produced")?;
139+
assert_eq!(3, batch.num_columns());
140+
assert_eq!(2, batch.num_rows());
141+
let schema = batch.schema();
142+
assert_eq!("string_col", schema.field(0).name());
143+
assert_eq!(&DataType::Binary, schema.field(0).data_type());
144+
let col = batch
145+
.column(0)
146+
.as_any()
147+
.downcast_ref::<BinaryArray>()
148+
.expect("column 0 not BinaryArray");
149+
assert_eq!("0".as_bytes(), col.value(0));
150+
assert_eq!("1".as_bytes(), col.value(1));
151+
assert_eq!("double_col", schema.field(1).name());
152+
assert_eq!(&DataType::Float64, schema.field(1).data_type());
153+
let col = batch
154+
.column(1)
155+
.as_any()
156+
.downcast_ref::<Float64Array>()
157+
.expect("column 1 not Float64Array");
158+
assert_eq!(0.0, col.value(0));
159+
assert_eq!(10.1, col.value(1));
160+
assert_eq!("bool_col", schema.field(2).name());
161+
assert_eq!(&DataType::Boolean, schema.field(2).data_type());
162+
let col = batch
163+
.column(2)
164+
.as_any()
165+
.downcast_ref::<BooleanArray>()
166+
.expect("column 2 not BooleanArray");
167+
assert!(col.value(0));
168+
assert!(!col.value(1));
169+
Ok(())
170+
}
171+
172+
fn build_projected_reader_schema(
173+
avro_json: &str,
174+
projected_field_names: &[&str],
175+
) -> DFResult<AvroSchema> {
176+
let mut schema_json: Value = serde_json::from_str(avro_json).map_err(|e| {
177+
DataFusionError::Execution(format!(
178+
"Failed to parse Avro schema JSON from metadata: {e}"
179+
))
180+
})?;
181+
let obj = schema_json.as_object_mut().ok_or_else(|| {
182+
DataFusionError::Execution(
183+
"Top-level Avro schema JSON is not an object".to_string(),
184+
)
185+
})?;
186+
let fields_val = obj.get_mut("fields").ok_or_else(|| {
187+
DataFusionError::Execution(
188+
"Top-level Avro schema JSON has no `fields` key".to_string(),
189+
)
190+
})?;
191+
let fields = fields_val.as_array_mut().ok_or_else(|| {
192+
DataFusionError::Execution(
193+
"Top-level Avro schema `fields` is not an array".to_string(),
194+
)
195+
})?;
196+
let mut by_name: HashMap<String, Value> = HashMap::new();
197+
for field in fields.iter() {
198+
if let Some(name) = field.get("name").and_then(|v| v.as_str()) {
199+
by_name.insert(name.to_string(), field.clone());
200+
}
201+
}
202+
let mut projected_fields = Vec::with_capacity(projected_field_names.len());
203+
for name in projected_field_names {
204+
let Some(field) = by_name.get(*name) else {
205+
return Err(DataFusionError::Execution(format!(
206+
"Projected field `{name}` not found in Avro writer schema"
207+
)));
208+
};
209+
projected_fields.push(field.clone());
210+
}
211+
*fields_val = Value::Array(projected_fields);
212+
let projected_json = serde_json::to_string(&schema_json).map_err(|e| {
213+
DataFusionError::Execution(format!(
214+
"Failed to serialize projected Avro schema JSON: {e}"
215+
))
216+
})?;
217+
Ok(AvroSchema::new(projected_json))
218+
}
44219
}

datafusion/datasource-avro/src/source.rs

Lines changed: 88 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -58,22 +58,99 @@ impl AvroSource {
5858
}
5959

6060
fn open<R: std::io::BufRead>(&self, reader: R) -> Result<Reader<R>> {
61-
let schema = self.table_schema.file_schema().as_ref(); // todo - avro metadata loading
62-
63-
let projected_schema = if let Some(projection) = &self.file_projection {
64-
&schema.project(projection)?
65-
} else {
66-
schema
67-
};
68-
69-
let avro_schema = AvroSchema::try_from(projected_schema)?;
70-
61+
// TODO: Once `ReaderBuilder::with_projection` is available, we should use it instead.
62+
// This should be an easy change. We'd simply need to:
63+
// 1. Use the full file schema to generate the reader `AvroSchema`.
64+
// 2. Pass `&self.file_projection` into `ReaderBuilder::with_projection`.
65+
// 3. Remove the `build_projected_reader_schema` methods.
7166
ReaderBuilder::new()
72-
.with_reader_schema(avro_schema) // Used for projection on read.
67+
.with_reader_schema(self.build_projected_reader_schema()?)
7368
.with_batch_size(self.batch_size.expect("Batch size must set before open"))
7469
.build(reader)
7570
.map_err(Into::into)
7671
}
72+
73+
fn build_projected_reader_schema(&self) -> Result<AvroSchema> {
74+
let file_schema = self.table_schema.file_schema().as_ref();
75+
// Fast path: no projection. If we have the original writer schema JSON
76+
// in metadata, just reuse it as-is without parsing.
77+
if self.file_projection.is_none() {
78+
return if let Some(avro_json) =
79+
file_schema.metadata().get(SCHEMA_METADATA_KEY)
80+
{
81+
Ok(AvroSchema::new(avro_json.clone()))
82+
} else {
83+
// Fall back to deriving Avro from the full Arrow file schema, should be ok
84+
// if not using projection.
85+
Ok(AvroSchema::try_from(file_schema)
86+
.map_err(Into::<DataFusionError>::into)?)
87+
};
88+
}
89+
// Use the writer Avro schema JSON tagged upstream to build a projected reader schema
90+
match file_schema.metadata().get(SCHEMA_METADATA_KEY) {
91+
Some(avro_json) => {
92+
let mut schema_json: Value =
93+
serde_json::from_str(avro_json).map_err(|e| {
94+
DataFusionError::Execution(format!(
95+
"Failed to parse Avro schema JSON from metadata: {e}"
96+
))
97+
})?;
98+
let obj = schema_json.as_object_mut().ok_or_else(|| {
99+
DataFusionError::Execution(
100+
"Top-level Avro schema JSON must be an object".to_string(),
101+
)
102+
})?;
103+
let fields_val = obj.get_mut("fields").ok_or_else(|| {
104+
DataFusionError::Execution(
105+
"Top-level Avro schema JSON must contain a `fields` array"
106+
.to_string(),
107+
)
108+
})?;
109+
let fields_arr = fields_val.as_array_mut().ok_or_else(|| {
110+
DataFusionError::Execution(
111+
"Top-level Avro schema `fields` must be an array".to_string(),
112+
)
113+
})?;
114+
// Move existing fields out so we can rebuild them in projected order.
115+
let original_fields = std::mem::take(fields_arr);
116+
let mut by_name: HashMap<String, Value> =
117+
HashMap::with_capacity(original_fields.len());
118+
for field in original_fields {
119+
if let Some(name) = field.get("name").and_then(|v| v.as_str()) {
120+
by_name.insert(name.to_string(), field);
121+
}
122+
}
123+
// Rebuild `fields` in the same order as the projected Arrow schema.
124+
let projection = self.file_projection.as_ref().ok_or_else(|| {
125+
DataFusionError::Internal("checked file_projection is Some above".to_string())
126+
})?;
127+
let projected_schema = file_schema.project(projection)?;
128+
let mut projected_fields =
129+
Vec::with_capacity(projected_schema.fields().len());
130+
for arrow_field in projected_schema.fields() {
131+
let name = arrow_field.name();
132+
let field = by_name.remove(name).ok_or_else(|| {
133+
DataFusionError::Execution(format!(
134+
"Projected field `{name}` not found in Avro writer schema"
135+
))
136+
})?;
137+
projected_fields.push(field);
138+
}
139+
*fields_val = Value::Array(projected_fields);
140+
let projected_json =
141+
serde_json::to_string(&schema_json).map_err(|e| {
142+
DataFusionError::Execution(format!(
143+
"Failed to serialize projected Avro schema JSON: {e}"
144+
))
145+
})?;
146+
Ok(AvroSchema::new(projected_json))
147+
}
148+
None => Err(DataFusionError::Execution(format!(
149+
"Avro schema metadata ({SCHEMA_METADATA_KEY}) is missing from file schema, but is required for projection"
150+
))),
151+
}
152+
}
153+
}
77154
}
78155

79156
impl FileSource for AvroSource {

0 commit comments

Comments
 (0)