Skip to main content

file_engine/query/
file_stream.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16
17use common_datasource::file_format::Format;
18use common_datasource::file_format::csv::CsvFormat;
19use common_datasource::file_format::orc::OrcSource;
20use common_datasource::file_format::parquet::DefaultParquetFileReaderFactory;
21use datafusion::common::ToDFSchema;
22use datafusion::config::CsvOptions;
23use datafusion::datasource::listing::PartitionedFile;
24use datafusion::datasource::object_store::ObjectStoreUrl;
25use datafusion::datasource::physical_plan::{
26    CsvSource, FileGroup, FileScanConfigBuilder, FileSource, FileStreamBuilder, JsonSource,
27    ParquetSource,
28};
29use datafusion::datasource::source::DataSourceExec;
30use datafusion::physical_expr::create_physical_expr;
31use datafusion::physical_expr::execution_props::ExecutionProps;
32use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
33use datafusion::physical_plan::{
34    ExecutionPlan, SendableRecordBatchStream as DfSendableRecordBatchStream,
35};
36use datafusion::prelude::SessionContext;
37use datafusion_expr::expr::Expr;
38use datafusion_expr::physical_planning_context::PhysicalPlanningContext;
39use datafusion_expr::utils::conjunction;
40use datatypes::schema::SchemaRef;
41use object_store::ObjectStore;
42use snafu::ResultExt;
43
44use crate::error::{self, Result};
45
46const DEFAULT_BATCH_SIZE: usize = 8192;
47
48fn build_record_batch_stream(
49    scan_plan_config: &ScanPlanConfig,
50    limit: Option<usize>,
51    file_source: Arc<dyn FileSource>,
52) -> Result<DfSendableRecordBatchStream> {
53    let files = scan_plan_config
54        .files
55        .iter()
56        .map(|filename| PartitionedFile::new(filename.clone(), 0))
57        .collect::<Vec<_>>();
58
59    let config =
60        FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), file_source.clone())
61            .with_projection_indices(scan_plan_config.projection.cloned())?
62            .with_limit(limit)
63            .with_file_group(FileGroup::new(files))
64            .build();
65
66    let store = Arc::new(object_store_opendal::OpendalStore::new(
67        scan_plan_config.store.clone(),
68    ));
69
70    let file_opener = config.file_source().create_file_opener(store, &config, 0)?;
71    let stream = FileStreamBuilder::new(&config)
72        .with_partition(0) // partition: hard-code
73        .with_file_opener(file_opener)
74        .with_metrics(&ExecutionPlanMetricsSet::new())
75        .build()
76        .context(error::BuildStreamSnafu)?;
77    Ok(Box::pin(stream))
78}
79
80fn new_csv_stream(
81    config: &ScanPlanConfig,
82    format: &CsvFormat,
83) -> Result<DfSendableRecordBatchStream> {
84    let file_schema = config.file_schema.arrow_schema().clone();
85
86    // push down limit only if there is no filter
87    let limit = config.filters.is_empty().then_some(config.limit).flatten();
88
89    let options = CsvOptions::default()
90        .with_has_header(format.has_header)
91        .with_delimiter(format.delimiter);
92    let csv_source = CsvSource::new(file_schema)
93        .with_csv_options(options)
94        .with_batch_size(DEFAULT_BATCH_SIZE);
95
96    build_record_batch_stream(config, limit, csv_source)
97}
98
99fn new_json_stream(config: &ScanPlanConfig) -> Result<DfSendableRecordBatchStream> {
100    let file_schema = config.file_schema.arrow_schema().clone();
101
102    // push down limit only if there is no filter
103    let limit = config.filters.is_empty().then_some(config.limit).flatten();
104
105    let file_source = JsonSource::new(file_schema).with_batch_size(DEFAULT_BATCH_SIZE);
106    build_record_batch_stream(config, limit, file_source)
107}
108
109fn new_parquet_stream_with_exec_plan(
110    config: &ScanPlanConfig,
111) -> Result<DfSendableRecordBatchStream> {
112    let file_schema = config.file_schema.arrow_schema().clone();
113    let ScanPlanConfig {
114        files,
115        projection,
116        limit,
117        filters,
118        store,
119        ..
120    } = config;
121
122    let file_group = FileGroup::new(
123        files
124            .iter()
125            .map(|filename| PartitionedFile::new(filename.clone(), 0))
126            .collect::<Vec<_>>(),
127    );
128
129    let mut parquet_source = ParquetSource::new(file_schema.clone())
130        .with_parquet_file_reader_factory(Arc::new(DefaultParquetFileReaderFactory::new(
131            store.clone(),
132        )));
133
134    // build predicate filter
135    let filters = filters.to_vec();
136    if let Some(expr) = conjunction(filters) {
137        let df_schema = file_schema
138            .clone()
139            .to_dfschema_ref()
140            .context(error::ParquetScanPlanSnafu)?;
141
142        let filters = create_physical_expr(
143            &expr,
144            &df_schema,
145            &ExecutionProps::new(),
146            &PhysicalPlanningContext::default(),
147        )
148        .context(error::ParquetScanPlanSnafu)?;
149        parquet_source = parquet_source.with_predicate(filters);
150    };
151
152    let file_scan_config =
153        FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), Arc::new(parquet_source))
154            .with_file_group(file_group)
155            .with_projection_indices(projection.cloned())?
156            .with_limit(*limit)
157            .build();
158
159    // TODO(ruihang): get this from upper layer
160    let task_ctx = SessionContext::default().task_ctx();
161
162    let parquet_exec = DataSourceExec::from_data_source(file_scan_config);
163    let stream = parquet_exec
164        .execute(0, task_ctx)
165        .context(error::ParquetScanPlanSnafu)?;
166
167    Ok(stream)
168}
169
170fn new_orc_stream(config: &ScanPlanConfig) -> Result<DfSendableRecordBatchStream> {
171    let file_schema = config.file_schema.arrow_schema().clone();
172
173    // push down limit only if there is no filter
174    let limit = config.filters.is_empty().then_some(config.limit).flatten();
175
176    let file_source = OrcSource::new(file_schema.into()).with_batch_size(DEFAULT_BATCH_SIZE);
177    build_record_batch_stream(config, limit, file_source)
178}
179
180#[derive(Debug, Clone)]
181pub struct ScanPlanConfig<'a> {
182    pub file_schema: SchemaRef,
183    pub files: &'a Vec<String>,
184    pub projection: Option<&'a Vec<usize>>,
185    pub filters: &'a [Expr],
186    pub limit: Option<usize>,
187    pub store: ObjectStore,
188}
189
190pub fn create_stream(
191    format: &Format,
192    config: &ScanPlanConfig,
193) -> Result<DfSendableRecordBatchStream> {
194    match format {
195        Format::Csv(format) => new_csv_stream(config, format),
196        Format::Json(_) => new_json_stream(config),
197        Format::Parquet(_) => new_parquet_stream_with_exec_plan(config),
198        Format::Orc(_) => new_orc_stream(config),
199    }
200}