Encode document rows as Arrow IPC bytes for columnar transport. Returns `None` if rows are empty or schema inference fails.
(
rows: &[(String, serde_json::Value)],
projection: &[String],
)
| 17 | /// |
| 18 | /// Returns `None` if rows are empty or schema inference fails. |
| 19 | pub fn encode_as_arrow_ipc( |
| 20 | rows: &[(String, serde_json::Value)], |
| 21 | projection: &[String], |
| 22 | ) -> Option<Vec<u8>> { |
| 23 | if rows.is_empty() { |
| 24 | return None; |
| 25 | } |
| 26 | |
| 27 | let first_obj = rows[0].1.as_object()?; |
| 28 | let field_names: Vec<&str> = if projection.is_empty() { |
| 29 | first_obj.keys().map(|k| k.as_str()).collect() |
| 30 | } else { |
| 31 | projection.iter().map(|s| s.as_str()).collect() |
| 32 | }; |
| 33 | |
| 34 | if field_names.is_empty() { |
| 35 | return None; |
| 36 | } |
| 37 | |
| 38 | let mut fields = vec![Field::new("id", DataType::Utf8, false)]; |
| 39 | for &name in &field_names { |
| 40 | let dt = first_obj |
| 41 | .get(name) |
| 42 | .map(infer_type) |
| 43 | .unwrap_or(DataType::Utf8); |
| 44 | fields.push(Field::new(name, dt, true)); |
| 45 | } |
| 46 | let schema = Arc::new(Schema::new(fields)); |
| 47 | |
| 48 | let mut ids: Vec<String> = Vec::with_capacity(rows.len()); |
| 49 | let mut builders: Vec<ColBuilder> = field_names |
| 50 | .iter() |
| 51 | .map(|&name| { |
| 52 | let dt = first_obj |
| 53 | .get(name) |
| 54 | .map(infer_type) |
| 55 | .unwrap_or(DataType::Utf8); |
| 56 | ColBuilder::new(dt, rows.len()) |
| 57 | }) |
| 58 | .collect(); |
| 59 | |
| 60 | for (doc_id, data) in rows { |
| 61 | ids.push(doc_id.clone()); |
| 62 | let obj = data.as_object(); |
| 63 | for (i, &name) in field_names.iter().enumerate() { |
| 64 | match obj.and_then(|o| o.get(name)) { |
| 65 | Some(v) => builders[i].push(v), |
| 66 | None => builders[i].push_null(), |
| 67 | } |
| 68 | } |
| 69 | } |
| 70 | |
| 71 | let mut arrays: Vec<ArrayRef> = vec![Arc::new(StringArray::from(ids))]; |
| 72 | for b in builders { |
| 73 | arrays.push(b.finish()); |
| 74 | } |
| 75 | |
| 76 | let batch = RecordBatch::try_new(schema.clone(), arrays).ok()?; |
nothing calls this directly
no test coverage detected