query/query_engine/
runtime.rs1use 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
26pub type QueryRuntimeProviderRef = Arc<dyn QueryRuntimeProvider>;
28
29#[derive(Clone, Copy)]
31#[non_exhaustive]
32pub struct QueryRuntimeContext<'a> {
33 pub query_options: &'a QueryOptions,
35 pub resolved_memory_pool_size: usize,
37}
38
39impl<'a> QueryRuntimeContext<'a> {
40 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
49pub trait QueryRuntimeProvider: Send + Sync + 'static {
51 fn configure_session_config(&self, _ctx: QueryRuntimeContext<'_>, _config: &mut SessionConfig) {
53 }
54
55 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#[derive(Debug, Default)]
67pub struct DefaultQueryRuntimeProvider;
68
69impl DefaultQueryRuntimeProvider {
70 pub fn runtime_env_builder(ctx: QueryRuntimeContext<'_>) -> RuntimeEnvBuilder {
72 let mut builder = RuntimeEnvBuilder::new();
73
74 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 }
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 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
129fn 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}