MCPcopy Create free account
hub / github.com/apache/datafusion / plan_to_csv

Function plan_to_csv

datafusion/datasource-csv/src/source.rs:437–494  ·  view source on GitHub ↗
(
    task_ctx: Arc<TaskContext>,
    plan: Arc<dyn ExecutionPlan>,
    path: impl AsRef<str>,
)

Source from the content-addressed store, hash-verified

435}
436
437pub async fn plan_to_csv(
438 task_ctx: Arc<TaskContext>,
439 plan: Arc<dyn ExecutionPlan>,
440 path: impl AsRef<str>,
441) -> Result<()> {
442 let path = path.as_ref();
443 let parsed = ListingTableUrl::parse(path)?;
444 let object_store_url = parsed.object_store();
445 let store = task_ctx.runtime_env().object_store(&object_store_url)?;
446 let writer_buffer_size = task_ctx
447 .session_config()
448 .options()
449 .execution
450 .objectstore_writer_buffer_size;
451 let mut join_set = JoinSet::new();
452 for i in 0..plan.output_partitioning().partition_count() {
453 let storeref = Arc::clone(&store);
454 let plan: Arc<dyn ExecutionPlan> = Arc::clone(&plan);
455 let filename = format!("{}/part-{i}.csv", parsed.prefix());
456 let file = object_store::path::Path::parse(filename)?;
457
458 let mut stream = plan.execute(i, Arc::clone(&task_ctx))?;
459 join_set.spawn(async move {
460 let mut buf_writer =
461 BufWriter::with_capacity(storeref, file.clone(), writer_buffer_size);
462 let mut buffer = Vec::with_capacity(1024);
463 //only write headers on first iteration
464 let mut write_headers = true;
465 while let Some(batch) = stream.next().await.transpose()? {
466 let mut writer = csv::WriterBuilder::new()
467 .with_header(write_headers)
468 .build(buffer);
469 writer.write(&batch)?;
470 buffer = writer.into_inner();
471 buf_writer.write_all(&buffer).await?;
472 buffer.clear();
473 //prevent writing headers more than once
474 write_headers = false;
475 }
476 buf_writer.shutdown().await.map_err(DataFusionError::from)
477 });
478 }
479
480 while let Some(result) = join_set.join_next().await {
481 match result {
482 Ok(res) => res?, // propagate DataFusion error
483 Err(e) => {
484 if e.is_panic() {
485 std::panic::resume_unwind(e.into_panic());
486 } else {
487 unreachable!();
488 }
489 }
490 }
491 }
492
493 Ok(())
494}

Callers 1

write_csvMethod · 0.85

Calls 15

newFunction · 0.85
session_configMethod · 0.80
partition_countMethod · 0.80
join_nextMethod · 0.80
as_refMethod · 0.45
object_storeMethod · 0.45
runtime_envMethod · 0.45
optionsMethod · 0.45
output_partitioningMethod · 0.45
executeMethod · 0.45
spawnMethod · 0.45
cloneMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…