Skip to main content

query/query_engine/
runtime.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 datafusion::error::Result as DfResult;
18use datafusion::execution::context::SessionConfig;
19use datafusion::execution::disk_manager::{DiskManagerBuilder, DiskManagerMode};
20use datafusion::execution::runtime_env::{RuntimeEnv, RuntimeEnvBuilder};
21use datafusion_common::config::SpillCompression;
22
23use crate::options::{QueryOptions, QuerySpillCompression, QuerySpillMode};
24use crate::query_engine::state::MetricsMemoryPool;
25
26/// Reference-counted query runtime provider.
27pub type QueryRuntimeProviderRef = Arc<dyn QueryRuntimeProvider>;
28
29/// Context for building query runtime components.
30#[derive(Clone, Copy)]
31#[non_exhaustive]
32pub struct QueryRuntimeContext<'a> {
33    /// Query options used by the query engine.
34    pub query_options: &'a QueryOptions,
35    /// Resolved memory pool size in bytes.
36    pub resolved_memory_pool_size: usize,
37}
38
39impl<'a> QueryRuntimeContext<'a> {
40    /// Creates a new query runtime context.
41    pub fn new(query_options: &'a QueryOptions, resolved_memory_pool_size: usize) -> Self {
42        Self {
43            query_options,
44            resolved_memory_pool_size,
45        }
46    }
47}
48
49/// Provides DataFusion session and runtime setup for the query engine.
50pub trait QueryRuntimeProvider: Send + Sync + 'static {
51    /// Configures the DataFusion session config before building the session state.
52    fn configure_session_config(&self, _ctx: QueryRuntimeContext<'_>, _config: &mut SessionConfig) {
53    }
54
55    /// Builds the DataFusion runtime environment.
56    fn build_runtime_env(
57        &self,
58        _ctx: QueryRuntimeContext<'_>,
59        builder: RuntimeEnvBuilder,
60    ) -> DfResult<Arc<RuntimeEnv>> {
61        builder.build().map(Arc::new)
62    }
63}
64
65/// Default query runtime provider.
66#[derive(Debug, Default)]
67pub struct DefaultQueryRuntimeProvider;
68
69impl DefaultQueryRuntimeProvider {
70    /// Creates a default DataFusion runtime environment builder.
71    pub fn runtime_env_builder(ctx: QueryRuntimeContext<'_>) -> RuntimeEnvBuilder {
72        let mut builder = RuntimeEnvBuilder::new();
73
74        // Attach the bounded metrics memory pool only when a limit is set
75        // (>0). When unbounded (0), keep the DataFusion default
76        // (UnboundedMemoryPool).
77        if ctx.resolved_memory_pool_size > 0 {
78            builder = builder.with_memory_pool(Arc::new(MetricsMemoryPool::new(
79                ctx.resolved_memory_pool_size,
80                ctx.query_options.experimental_memory_pool_policy,
81            )));
82        }
83
84        match ctx.query_options.experimental_spill_mode {
85            QuerySpillMode::Default => {
86                // No custom disk manager; preserve DataFusion default OS temp directory.
87            }
88            QuerySpillMode::Custom => {
89                let mut dm_builder = DiskManagerBuilder::default();
90                if let Some(ref path) = ctx.query_options.experimental_spill_path {
91                    dm_builder =
92                        dm_builder.with_mode(DiskManagerMode::Directories(vec![path.clone()]));
93                }
94                let max_temp_directory_size =
95                    ctx.query_options.experimental_spill_max_temp_directory_size;
96                dm_builder =
97                    dm_builder.with_max_temp_directory_size(max_temp_directory_size.as_bytes());
98                common_telemetry::info!(
99                    "Configured custom query spill: path={:?}, max_temp_directory_size={}, compression={:?}",
100                    ctx.query_options.experimental_spill_path,
101                    max_temp_directory_size,
102                    ctx.query_options.experimental_spill_compression,
103                );
104                builder = builder.with_disk_manager_builder(dm_builder);
105            }
106            QuerySpillMode::Disabled => {
107                let dm_builder = DiskManagerBuilder::default().with_mode(DiskManagerMode::Disabled);
108                builder = builder.with_disk_manager_builder(dm_builder);
109            }
110        }
111
112        builder
113    }
114}
115
116impl QueryRuntimeProvider for DefaultQueryRuntimeProvider {
117    fn configure_session_config(&self, ctx: QueryRuntimeContext<'_>, config: &mut SessionConfig) {
118        // Set spill compression on the session config only when spill mode is
119        // Custom. In Default/Disabled modes, DataFusion's own default
120        // (Uncompressed) is preserved—setting compression when spill is not
121        // explicitly configured would be misleading.
122        if ctx.query_options.experimental_spill_mode == QuerySpillMode::Custom {
123            config.options_mut().execution.spill_compression =
124                spill_compression_from_options(ctx.query_options.experimental_spill_compression);
125        }
126    }
127}
128
129/// Map [`QuerySpillCompression`] to DataFusion's [`SpillCompression`].
130///
131/// This conversion is intentionally not a `From` impl because the
132/// semantics depend on the spill mode; callers should only invoke
133/// this when `experimental_spill_mode == Custom`.
134fn spill_compression_from_options(comp: QuerySpillCompression) -> SpillCompression {
135    match comp {
136        QuerySpillCompression::Uncompressed => SpillCompression::Uncompressed,
137        QuerySpillCompression::Lz4Frame => SpillCompression::Lz4Frame,
138        QuerySpillCompression::Zstd => SpillCompression::Zstd,
139    }
140}