Skip to main content

datafusion_datasource_arrow/
source.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Execution plan for reading Arrow IPC files
19//!
20//! # Naming Note
21//!
22//! The naming in this module can be confusing:
23//! - `ArrowFileOpener` handles the Arrow IPC **file format**
24//!   (with footer, supports parallel reading)
25//! - `ArrowStreamFileOpener` handles the Arrow IPC **stream format**
26//!   (without footer, sequential only)
27//! - `ArrowSource` is the unified `FileSource` implementation that uses either opener
28//!   depending on the format specified at construction
29//!
30//! Despite the name "ArrowStreamFileOpener", it still reads from files - the "Stream"
31//! refers to the Arrow IPC stream format, not streaming I/O. Both formats can be stored
32//! in files on disk or object storage.
33
34use 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/// Enum indicating which Arrow IPC format to use
58#[derive(Clone, Copy, Debug)]
59enum ArrowFormat {
60    /// Arrow IPC file format (with footer, supports parallel reading)
61    File,
62    /// Arrow IPC stream format (without footer, sequential only)
63    Stream,
64}
65
66/// `FileOpener` for Arrow IPC stream format. Supports only sequential reading.
67pub(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
111/// `FileOpener` for Arrow IPC file format. Supports range-based parallel reading.
112pub(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                    // range is not none, the file maybe split into multiple parts to scan in parallel
152                    // get footer_len firstly
153                    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                    // read footer according to footer_len
165                    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                    // build decoder according to footer & projection
180                    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                    // filter recordbatches according to range
208                    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/// `FileSource` for both Arrow IPC file and stream formats
258#[derive(Clone)]
259pub struct ArrowSource {
260    format: ArrowFormat,
261    metrics: ExecutionPlanMetricsSet,
262    projection: SplitProjection,
263    table_schema: TableSchema,
264}
265
266impl ArrowSource {
267    /// Creates an [`ArrowSource`] for file format
268    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    /// Creates an [`ArrowSource`] for stream format
279    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                // The Arrow IPC stream format doesn't support range-based parallel reading
341                // because it lacks a footer with the information that would be needed to
342                // make range-based parallel reading practical. Without the data in the
343                // footer you would either need to read the the entire file and record the
344                // offsets of the record batches and dictionaries, essentially recreating
345                // the footer's contents, or else each partition would need to read the
346                // entire file up to the correct offset which is a lot of duplicate I/O.
347                // We're opting to avoid that entirely by only acting on a single partition
348                // and reading sequentially.
349                Ok(None)
350            }
351            ArrowFormat::File => {
352                // Use the default trait implementation logic for file format
353                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
397/// `FileOpener` wrapper for both Arrow IPC file and stream formats
398pub 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    /// Creates a new [`ArrowOpener`]
410    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]), // just the first column
629        };
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}