(
task_ctx: Arc<TaskContext>,
plan: Arc<dyn ExecutionPlan>,
path: impl AsRef<str>,
)
| 435 | } |
| 436 | |
| 437 | pub 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 | } |
no test coverage detected
searching dependent graphs…