MCPcopy Create free account
hub / github.com/Pometry/Raphtory / process_csv_paths_df

Function process_csv_paths_df

raphtory/src/python/graph/io/arrow_loaders.rs:776–839  ·  view source on GitHub ↗
(
    paths: &'a [PathBuf],
    col_names: &'a [&'a str],
    csv_options: Option<&'a CsvReadOptions>,
    schema: Option<Arc<HashMap<String, PropType>>>,
)

Source from the content-addressed store, hash-verified

774}
775
776fn 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 }),

Calls 15

ErrEnum · 0.85
build_csv_readerFunction · 0.85
cast_columnsFunction · 0.85
to_stringMethod · 0.80
enumerateMethod · 0.80
is_emptyMethod · 0.45
collectMethod · 0.45
mapMethod · 0.45
iterMethod · 0.45
cloneMethod · 0.45
lenMethod · 0.45
findMethod · 0.45

Tested by

no test coverage detected