Skip to main content

flow/
batching_mode.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
15//! Run flow as batching mode which is time-window-aware normal query triggered when new data arrives
16
17use std::time::Duration;
18
19use common_grpc::channel_manager::ClientTlsOption;
20use serde::{Deserialize, Serialize};
21use session::ReadPreference;
22
23pub(crate) mod batching_execution;
24pub(crate) mod checkpoint;
25pub(crate) mod engine;
26mod eval_schedule;
27pub(crate) mod frontend_client;
28pub(crate) mod state;
29mod table_creator;
30pub(crate) mod task;
31pub(crate) mod time_window;
32pub(crate) mod utils;
33
34#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
35pub struct BatchingModeOptions {
36    /// The default batching engine query timeout is 10 minutes
37    #[serde(with = "humantime_serde")]
38    pub query_timeout: Duration,
39    /// will output a warn log for any query that runs for more that this threshold
40    #[serde(with = "humantime_serde")]
41    pub slow_query_threshold: Duration,
42    /// The minimum duration between two queries execution by batching mode task
43    #[serde(with = "humantime_serde")]
44    pub experimental_min_refresh_duration: Duration,
45    /// The gRPC connection timeout
46    #[serde(with = "humantime_serde")]
47    pub grpc_conn_timeout: Duration,
48    #[serde(with = "humantime_serde")]
49    pub experimental_flight_do_get_timeout: Duration,
50    /// The gRPC max retry number
51    pub experimental_grpc_max_retries: u32,
52    /// Flow wait for available frontend timeout,
53    /// if failed to find available frontend after frontend_scan_timeout elapsed, return error
54    /// which prevent flownode from starting
55    #[serde(with = "humantime_serde")]
56    pub experimental_frontend_scan_timeout: Duration,
57    /// Maximum number of filters allowed in a single query
58    pub experimental_max_filter_num_per_query: usize,
59    /// Time window merge distance
60    pub experimental_time_window_merge_threshold: usize,
61    /// Whether to enable experimental flow incremental source reads.
62    ///
63    /// When disabled, batching flows always execute full-snapshot queries.
64    pub experimental_enable_incremental_read: bool,
65    /// Read preference of the Frontend client.
66    pub read_preference: ReadPreference,
67    /// TLS option for client connections to frontends.
68    pub frontend_tls: Option<ClientTlsOption>,
69}
70
71impl Default for BatchingModeOptions {
72    fn default() -> Self {
73        Self {
74            query_timeout: Duration::from_secs(10 * 60),
75            slow_query_threshold: Duration::from_secs(60),
76            experimental_min_refresh_duration: Duration::new(5, 0),
77            grpc_conn_timeout: Duration::from_secs(5),
78            experimental_flight_do_get_timeout: Duration::from_secs(10),
79            experimental_grpc_max_retries: 3,
80            experimental_frontend_scan_timeout: Duration::from_secs(30),
81            experimental_max_filter_num_per_query: 20,
82            experimental_time_window_merge_threshold: 3,
83            experimental_enable_incremental_read: false,
84            read_preference: Default::default(),
85            frontend_tls: None,
86        }
87    }
88}