(
paths: &'a [PathBuf],
col_names: &'a [&'a str],
csv_options: Option<&'a CsvReadOptions>,
schema: Option<Arc<HashMap<String, PropType>>>,
)
| 774 | } |
| 775 | |
| 776 | fn process_csv_paths_df<'a>( |
| 777 | paths: &'a [PathBuf], |
| 778 | col_names: &'a [&'a str], |
| 779 | csv_options: Option<&'a CsvReadOptions>, |
| 780 | schema: Option<Arc<HashMap<String, PropType>>>, |
| 781 | ) -> Result<DFView<impl Iterator<Item = Result<DFChunk, GraphError>> + Send + 'a>, GraphError> { |
| 782 | if paths.is_empty() { |
| 783 | return Err(GraphError::LoadFailure( |
| 784 | "No CSV files found at the provided path".to_string(), |
| 785 | )); |
| 786 | } |
| 787 | // BoxedLIter couldn't be used because it has Send + Sync bound |
| 788 | // type ChunkIter<'b> = Box<dyn Iterator<Item = Result<DFChunk, GraphError>> + 'b>; |
| 789 | |
| 790 | let names = col_names.iter().map(|&name| name.to_string()).collect(); |
| 791 | let chunks = paths.iter().flat_map(move |path| { |
| 792 | let schema = schema.clone(); |
| 793 | let csv_reader = match build_csv_reader(path.as_path(), csv_options) { |
| 794 | Ok(r) => r, |
| 795 | Err(e) => return Iter3::I(iter::once(Err(e))), |
| 796 | }; |
| 797 | let mut indices = Vec::with_capacity(col_names.len()); |
| 798 | for required_col in col_names { |
| 799 | if let Some((idx, _)) = csv_reader |
| 800 | .schema() |
| 801 | .fields() |
| 802 | .iter() |
| 803 | .enumerate() |
| 804 | .find(|(_, f)| f.name() == required_col) |
| 805 | { |
| 806 | indices.push(idx); |
| 807 | } else { |
| 808 | return Iter3::J(iter::once(Err(GraphError::ColumnDoesNotExist( |
| 809 | required_col.to_string(), |
| 810 | )))); |
| 811 | } |
| 812 | } |
| 813 | Iter3::K( |
| 814 | csv_reader |
| 815 | .into_iter() |
| 816 | .map(move |batch_res| match batch_res { |
| 817 | Ok(batch) => { |
| 818 | let casted_batch = if let Some(schema) = schema.as_deref() { |
| 819 | cast_columns(batch, schema)? |
| 820 | } else { |
| 821 | batch |
| 822 | }; |
| 823 | let arrays = indices |
| 824 | .iter() |
| 825 | .map(|&idx| casted_batch.column(idx).clone()) |
| 826 | .collect::<Vec<_>>(); |
| 827 | Ok(DFChunk::new(arrays)) |
| 828 | } |
| 829 | Err(e) => Err(GraphError::LoadFailure(format!( |
| 830 | "Arrow CSV error while reading a batch from '{}': {e}", |
| 831 | path.display() |
| 832 | ))), |
| 833 | }), |
no test coverage detected