datafusion_datasource_arrow/
source.rs1use std::io::Cursor;
35use std::sync::Arc;
36
37use datafusion_datasource::{TableSchema, as_file_source};
38
39use arrow::buffer::Buffer;
40use arrow::ipc::reader::{FileDecoder, FileReader, StreamReader};
41use datafusion_common::error::Result;
42use datafusion_common::exec_datafusion_err;
43use datafusion_datasource::PartitionedFile;
44use datafusion_datasource::file::FileSource;
45use datafusion_datasource::file_scan_config::FileScanConfig;
46use datafusion_datasource::projection::{ProjectionOpener, SplitProjection};
47use datafusion_physical_expr_common::sort_expr::LexOrdering;
48use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet;
49use datafusion_physical_plan::projection::ProjectionExprs;
50
51use datafusion_datasource::file_stream::FileOpenFuture;
52use datafusion_datasource::file_stream::FileOpener;
53use futures::StreamExt;
54use itertools::Itertools;
55use object_store::{GetOptions, GetRange, GetResultPayload, ObjectStore, ObjectStoreExt};
56
57#[derive(Clone, Copy, Debug)]
59enum ArrowFormat {
60 File,
62 Stream,
64}
65
66pub(crate) struct ArrowStreamFileOpener {
68 object_store: Arc<dyn ObjectStore>,
69 projection: Option<Vec<usize>>,
70}
71
72impl FileOpener for ArrowStreamFileOpener {
73 fn open(&self, partitioned_file: PartitionedFile) -> Result<FileOpenFuture> {
74 if partitioned_file.range.is_some() {
75 return Err(exec_datafusion_err!(
76 "ArrowStreamFileOpener does not support range-based reading"
77 ));
78 }
79 let object_store = Arc::clone(&self.object_store);
80 let projection = self.projection.clone();
81
82 Ok(Box::pin(async move {
83 let r = object_store
84 .get(&partitioned_file.object_meta.location)
85 .await?;
86
87 let stream = match r.payload {
88 #[cfg(not(target_arch = "wasm32"))]
89 GetResultPayload::File(file, _) => futures::stream::iter(
90 StreamReader::try_new(file.try_clone()?, projection.clone())?,
91 )
92 .map(|r| r.map_err(Into::into))
93 .boxed(),
94 GetResultPayload::Stream(_) => {
95 let bytes = r.bytes().await?;
96 let cursor = Cursor::new(bytes);
97 futures::stream::iter(StreamReader::try_new(
98 cursor,
99 projection.clone(),
100 )?)
101 .map(|r| r.map_err(Into::into))
102 .boxed()
103 }
104 };
105
106 Ok(stream)
107 }))
108 }
109}
110
111pub(crate) struct ArrowFileOpener {
113 object_store: Arc<dyn ObjectStore>,
114 projection: Option<Vec<usize>>,
115}
116
117impl FileOpener for ArrowFileOpener {
118 fn open(&self, partitioned_file: PartitionedFile) -> Result<FileOpenFuture> {
119 let object_store = Arc::clone(&self.object_store);
120 let projection = self.projection.clone();
121
122 Ok(Box::pin(async move {
123 let range = partitioned_file.range.clone();
124 match range {
125 None => {
126 let r = object_store
127 .get(&partitioned_file.object_meta.location)
128 .await?;
129 let stream = match r.payload {
130 #[cfg(not(target_arch = "wasm32"))]
131 GetResultPayload::File(file, _) => futures::stream::iter(
132 FileReader::try_new(file.try_clone()?, projection.clone())?,
133 )
134 .map(|r| r.map_err(Into::into))
135 .boxed(),
136 GetResultPayload::Stream(_) => {
137 let bytes = r.bytes().await?;
138 let cursor = Cursor::new(bytes);
139 futures::stream::iter(FileReader::try_new(
140 cursor,
141 projection.clone(),
142 )?)
143 .map(|r| r.map_err(Into::into))
144 .boxed()
145 }
146 };
147
148 Ok(stream)
149 }
150 Some(range) => {
151 let get_option = GetOptions {
154 range: Some(GetRange::Suffix(10)),
155 ..Default::default()
156 };
157 let get_result = object_store
158 .get_opts(&partitioned_file.object_meta.location, get_option)
159 .await?;
160 let footer_len_buf = get_result.bytes().await?;
161 let footer_len = arrow_ipc::reader::read_footer_length(
162 footer_len_buf[..].try_into().unwrap(),
163 )?;
164 let get_option = GetOptions {
166 range: Some(GetRange::Suffix(10 + (footer_len as u64))),
167 ..Default::default()
168 };
169 let get_result = object_store
170 .get_opts(&partitioned_file.object_meta.location, get_option)
171 .await?;
172 let footer_buf = get_result.bytes().await?;
173 let footer = arrow_ipc::root_as_footer(
174 footer_buf[..footer_len].try_into().unwrap(),
175 )
176 .map_err(|err| {
177 exec_datafusion_err!("Unable to get root as footer: {err:?}")
178 })?;
179 let schema =
181 arrow_ipc::convert::fb_to_schema(footer.schema().unwrap());
182 let mut decoder = FileDecoder::new(schema.into(), footer.version());
183 if let Some(projection) = projection {
184 decoder = decoder.with_projection(projection);
185 }
186 let dict_ranges = footer
187 .dictionaries()
188 .iter()
189 .flatten()
190 .map(|block| {
191 let block_len =
192 block.bodyLength() as u64 + block.metaDataLength() as u64;
193 let block_offset = block.offset() as u64;
194 block_offset..block_offset + block_len
195 })
196 .collect_vec();
197 let dict_results = object_store
198 .get_ranges(&partitioned_file.object_meta.location, &dict_ranges)
199 .await?;
200 for (dict_block, dict_result) in
201 footer.dictionaries().iter().flatten().zip(dict_results)
202 {
203 decoder
204 .read_dictionary(dict_block, &Buffer::from(dict_result))?;
205 }
206
207 let recordbatches = footer
209 .recordBatches()
210 .iter()
211 .flatten()
212 .filter(|block| {
213 let block_offset = block.offset() as u64;
214 block_offset >= range.start as u64
215 && block_offset < range.end as u64
216 })
217 .copied()
218 .collect_vec();
219
220 let recordbatch_ranges = recordbatches
221 .iter()
222 .map(|block| {
223 let block_len =
224 block.bodyLength() as u64 + block.metaDataLength() as u64;
225 let block_offset = block.offset() as u64;
226 block_offset..block_offset + block_len
227 })
228 .collect_vec();
229
230 let recordbatch_results = object_store
231 .get_ranges(
232 &partitioned_file.object_meta.location,
233 &recordbatch_ranges,
234 )
235 .await?;
236
237 let stream = futures::stream::iter(
238 recordbatches
239 .into_iter()
240 .zip(recordbatch_results)
241 .filter_map(move |(block, data)| {
242 decoder
243 .read_record_batch(&block, &Buffer::from(data))
244 .transpose()
245 }),
246 )
247 .map(|r| r.map_err(Into::into))
248 .boxed();
249
250 Ok(stream)
251 }
252 }
253 }))
254 }
255}
256
257#[derive(Clone)]
259pub struct ArrowSource {
260 format: ArrowFormat,
261 metrics: ExecutionPlanMetricsSet,
262 projection: SplitProjection,
263 table_schema: TableSchema,
264}
265
266impl ArrowSource {
267 pub fn new_file_source(table_schema: impl Into<TableSchema>) -> Self {
269 let table_schema = table_schema.into();
270 Self {
271 format: ArrowFormat::File,
272 metrics: ExecutionPlanMetricsSet::new(),
273 projection: SplitProjection::unprojected(&table_schema),
274 table_schema,
275 }
276 }
277
278 pub fn new_stream_file_source(table_schema: impl Into<TableSchema>) -> Self {
280 let table_schema = table_schema.into();
281 Self {
282 format: ArrowFormat::Stream,
283 metrics: ExecutionPlanMetricsSet::new(),
284 projection: SplitProjection::unprojected(&table_schema),
285 table_schema,
286 }
287 }
288}
289
290impl FileSource for ArrowSource {
291 fn create_file_opener(
292 &self,
293 object_store: Arc<dyn ObjectStore>,
294 _base_config: &FileScanConfig,
295 _partition: usize,
296 ) -> Result<Arc<dyn FileOpener>> {
297 let split_projection = self.projection.clone();
298
299 let opener: Arc<dyn FileOpener> = match self.format {
300 ArrowFormat::File => Arc::new(ArrowFileOpener {
301 object_store,
302 projection: Some(split_projection.file_indices.clone()),
303 }),
304 ArrowFormat::Stream => Arc::new(ArrowStreamFileOpener {
305 object_store,
306 projection: Some(split_projection.file_indices.clone()),
307 }),
308 };
309 ProjectionOpener::try_new(
310 split_projection,
311 opener,
312 self.table_schema.file_schema(),
313 )
314 }
315
316 fn with_batch_size(&self, _batch_size: usize) -> Arc<dyn FileSource> {
317 Arc::new(Self { ..self.clone() })
318 }
319
320 fn metrics(&self) -> &ExecutionPlanMetricsSet {
321 &self.metrics
322 }
323
324 fn file_type(&self) -> &str {
325 match self.format {
326 ArrowFormat::File => "arrow",
327 ArrowFormat::Stream => "arrow_stream",
328 }
329 }
330
331 fn repartitioned(
332 &self,
333 target_partitions: usize,
334 repartition_file_min_size: usize,
335 output_ordering: Option<LexOrdering>,
336 config: &FileScanConfig,
337 ) -> Result<Option<FileScanConfig>> {
338 match self.format {
339 ArrowFormat::Stream => {
340 Ok(None)
350 }
351 ArrowFormat::File => {
352 use datafusion_datasource::file_groups::FileGroupPartitioner;
354
355 if config.file_compression_type.is_compressed() {
356 return Ok(None);
357 }
358
359 let repartitioned_file_groups_option = FileGroupPartitioner::new()
360 .with_target_partitions(target_partitions)
361 .with_repartition_file_min_size(repartition_file_min_size)
362 .with_preserve_order_within_groups(output_ordering.is_some())
363 .repartition_file_groups(&config.file_groups);
364
365 if let Some(repartitioned_file_groups) = repartitioned_file_groups_option
366 {
367 let mut source = config.clone();
368 source.file_groups = repartitioned_file_groups;
369 return Ok(Some(source));
370 }
371 Ok(None)
372 }
373 }
374 }
375
376 fn table_schema(&self) -> &TableSchema {
377 &self.table_schema
378 }
379
380 fn try_pushdown_projection(
381 &self,
382 projection: &ProjectionExprs,
383 ) -> Result<Option<Arc<dyn FileSource>>> {
384 let mut source = self.clone();
385 source.projection = SplitProjection::new(
386 self.table_schema().file_schema(),
387 &source.projection.source.try_merge(projection)?,
388 );
389 Ok(Some(Arc::new(source)))
390 }
391
392 fn projection(&self) -> Option<&ProjectionExprs> {
393 Some(&self.projection.source)
394 }
395}
396
397pub struct ArrowOpener {
399 pub inner: Arc<dyn FileOpener>,
400}
401
402impl FileOpener for ArrowOpener {
403 fn open(&self, partitioned_file: PartitionedFile) -> Result<FileOpenFuture> {
404 self.inner.open(partitioned_file)
405 }
406}
407
408impl ArrowOpener {
409 pub fn new(inner: Arc<dyn FileOpener>) -> Self {
411 Self { inner }
412 }
413
414 pub fn new_file_opener(
415 object_store: Arc<dyn ObjectStore>,
416 projection: Option<Vec<usize>>,
417 ) -> Self {
418 Self {
419 inner: Arc::new(ArrowFileOpener {
420 object_store,
421 projection,
422 }),
423 }
424 }
425
426 pub fn new_stream_file_opener(
427 object_store: Arc<dyn ObjectStore>,
428 projection: Option<Vec<usize>>,
429 ) -> Self {
430 Self {
431 inner: Arc::new(ArrowStreamFileOpener {
432 object_store,
433 projection,
434 }),
435 }
436 }
437}
438
439impl From<ArrowSource> for Arc<dyn FileSource> {
440 fn from(source: ArrowSource) -> Self {
441 as_file_source(source)
442 }
443}
444
445#[cfg(test)]
446mod tests {
447 use std::{fs::File, io::Read};
448
449 use arrow::datatypes::{DataType, Field, Schema};
450 use arrow_ipc::reader::{FileReader, StreamReader};
451 use bytes::Bytes;
452 use datafusion_datasource::file_scan_config::FileScanConfigBuilder;
453 use datafusion_execution::object_store::ObjectStoreUrl;
454 use object_store::memory::InMemory;
455
456 use super::*;
457
458 #[tokio::test]
459 async fn test_file_opener_without_ranges() -> Result<()> {
460 for filename in ["example.arrow", "example_stream.arrow"] {
461 let path = format!("tests/data/{filename}");
462 let path_str = path.as_str();
463 let mut file = File::open(path_str)?;
464 let file_size = file.metadata()?.len();
465
466 let mut buffer = Vec::new();
467 file.read_to_end(&mut buffer)?;
468 let bytes = Bytes::from(buffer);
469
470 let object_store = Arc::new(InMemory::new());
471 let partitioned_file = PartitionedFile::new(filename, file_size);
472 object_store
473 .put(&partitioned_file.object_meta.location, bytes.into())
474 .await?;
475
476 let schema = match FileReader::try_new(File::open(path_str)?, None) {
477 Ok(reader) => reader.schema(),
478 Err(_) => StreamReader::try_new(File::open(path_str)?, None)?.schema(),
479 };
480
481 let source: Arc<dyn FileSource> = if filename.contains("stream") {
482 Arc::new(ArrowSource::new_stream_file_source(schema))
483 } else {
484 Arc::new(ArrowSource::new_file_source(schema))
485 };
486
487 let scan_config = FileScanConfigBuilder::new(
488 ObjectStoreUrl::local_filesystem(),
489 source.clone(),
490 )
491 .build();
492
493 let file_opener = source.create_file_opener(object_store, &scan_config, 0)?;
494 let mut stream = file_opener.open(partitioned_file)?.await?;
495
496 assert!(stream.next().await.is_some());
497 }
498
499 Ok(())
500 }
501
502 #[tokio::test]
503 async fn test_file_opener_with_ranges() -> Result<()> {
504 let filename = "example.arrow";
505 let path = format!("tests/data/{filename}");
506 let path_str = path.as_str();
507 let mut file = File::open(path_str)?;
508 let file_size = file.metadata()?.len();
509
510 let mut buffer = Vec::new();
511 file.read_to_end(&mut buffer)?;
512 let bytes = Bytes::from(buffer);
513
514 let object_store = Arc::new(InMemory::new());
515 let partitioned_file = PartitionedFile::new_with_range(
516 filename.into(),
517 file_size,
518 0,
519 (file_size - 1) as i64,
520 );
521 object_store
522 .put(&partitioned_file.object_meta.location, bytes.into())
523 .await?;
524
525 let schema = FileReader::try_new(File::open(path_str)?, None)?.schema();
526
527 let source = Arc::new(ArrowSource::new_file_source(schema));
528
529 let scan_config = FileScanConfigBuilder::new(
530 ObjectStoreUrl::local_filesystem(),
531 source.clone(),
532 )
533 .build();
534
535 let file_opener = source.create_file_opener(object_store, &scan_config, 0)?;
536 let mut stream = file_opener.open(partitioned_file)?.await?;
537
538 assert!(stream.next().await.is_some());
539
540 Ok(())
541 }
542
543 #[tokio::test]
544 async fn test_stream_opener_errors_with_ranges() -> Result<()> {
545 let filename = "example_stream.arrow";
546 let path = format!("tests/data/{filename}");
547 let path_str = path.as_str();
548 let mut file = File::open(path_str)?;
549 let file_size = file.metadata()?.len();
550
551 let mut buffer = Vec::new();
552 file.read_to_end(&mut buffer)?;
553 let bytes = Bytes::from(buffer);
554
555 let object_store = Arc::new(InMemory::new());
556 let partitioned_file = PartitionedFile::new_with_range(
557 filename.into(),
558 file_size,
559 0,
560 (file_size - 1) as i64,
561 );
562 object_store
563 .put(&partitioned_file.object_meta.location, bytes.into())
564 .await?;
565
566 let schema = StreamReader::try_new(File::open(path_str)?, None)?.schema();
567
568 let source = Arc::new(ArrowSource::new_stream_file_source(schema));
569
570 let scan_config = FileScanConfigBuilder::new(
571 ObjectStoreUrl::local_filesystem(),
572 source.clone(),
573 )
574 .build();
575
576 let file_opener = source.create_file_opener(object_store, &scan_config, 0)?;
577 let result = file_opener.open(partitioned_file);
578 assert!(result.is_err());
579
580 Ok(())
581 }
582
583 #[tokio::test]
584 async fn test_arrow_stream_repartitioning_not_supported() -> Result<()> {
585 let schema =
586 Arc::new(Schema::new(vec![Field::new("f0", DataType::Int64, false)]));
587 let source = ArrowSource::new_stream_file_source(schema);
588
589 let config = FileScanConfigBuilder::new(
590 ObjectStoreUrl::local_filesystem(),
591 Arc::new(source.clone()) as Arc<dyn FileSource>,
592 )
593 .build();
594
595 for target_partitions in [2, 4, 8, 16] {
596 let result =
597 source.repartitioned(target_partitions, 1024 * 1024, None, &config)?;
598
599 assert!(
600 result.is_none(),
601 "Stream format should not support repartitioning with {target_partitions} partitions",
602 );
603 }
604
605 Ok(())
606 }
607
608 #[tokio::test]
609 async fn test_stream_opener_with_projection() -> Result<()> {
610 let filename = "example_stream.arrow";
611 let path = format!("tests/data/{filename}");
612 let path_str = path.as_str();
613 let mut file = File::open(path_str)?;
614 let file_size = file.metadata()?.len();
615
616 let mut buffer = Vec::new();
617 file.read_to_end(&mut buffer)?;
618 let bytes = Bytes::from(buffer);
619
620 let object_store = Arc::new(InMemory::new());
621 let partitioned_file = PartitionedFile::new(filename, file_size);
622 object_store
623 .put(&partitioned_file.object_meta.location, bytes.into())
624 .await?;
625
626 let opener = ArrowStreamFileOpener {
627 object_store,
628 projection: Some(vec![0]), };
630
631 let mut stream = opener.open(partitioned_file)?.await?;
632
633 if let Some(batch) = stream.next().await {
634 let batch = batch?;
635 assert_eq!(
636 batch.num_columns(),
637 1,
638 "Projection should result in 1 column"
639 );
640 } else {
641 panic!("Expected at least one batch");
642 }
643
644 Ok(())
645 }
646}