Skip to main content

flow/batching_mode/task/
inc.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_error::ext::BoxedError;
18use common_telemetry::debug;
19use common_telemetry::tracing::warn;
20use datafusion_expr::{DmlStatement, LogicalPlan};
21use query::QueryEngineRef;
22use query::options::{
23    FLOW_INCREMENTAL_AFTER_SEQS, FLOW_INCREMENTAL_MODE, FLOW_INCREMENTAL_MODE_MEMTABLE_ONLY,
24    FLOW_INCREMENTAL_MODE_SEQUENCE_RANGE, FLOW_SINK_TABLE_ID,
25};
26use snafu::ResultExt;
27use store_api::mito_engine_options::PRESERVE_ROW_SEQUENCE;
28use table::metadata::TableId;
29
30use crate::Error;
31use crate::batching_mode::state::CheckpointMode;
32use crate::batching_mode::table_creator::QueryType;
33use crate::batching_mode::task::BatchingTask;
34use crate::batching_mode::utils::{
35    analyze_incremental_aggregate_plan, get_table_info_df_schema,
36    rewrite_incremental_aggregate_with_sink_merge,
37};
38use crate::error::{ExternalSnafu, UnexpectedSnafu};
39
40impl BatchingTask {
41    async fn sink_table_id(&self) -> Result<TableId, Error> {
42        let table = self
43            .config
44            .catalog_manager
45            .table(
46                &self.config.sink_table_name[0],
47                &self.config.sink_table_name[1],
48                &self.config.sink_table_name[2],
49                None,
50            )
51            .await
52            .map_err(BoxedError::new)
53            .context(ExternalSnafu)?
54            .ok_or_else(|| {
55                UnexpectedSnafu {
56                    reason: format!(
57                        "Flow {} cannot build incremental extensions because sink table {:?} was not found",
58                        self.config.flow_id, self.config.sink_table_name
59                    ),
60                }
61                .build()
62            })?;
63        Ok(table.table_info().table_id())
64    }
65
66    /// Whether every source table provably supports exact sequence-range
67    /// scans: it must be the canonical mito engine and declare the
68    /// `preserve_row_sequence` capability (a mito region option enforced to
69    /// require append-only mode). When the capability is unknown — a source
70    /// table that cannot be resolved, is not the mito engine, or lacks the
71    /// option — this returns `false` so the caller keeps the historical
72    /// `memtable_only` mode instead of upgrading.
73    pub async fn sequence_range_capable(&self) -> Result<bool, Error> {
74        for name in &self.config.source_table_names {
75            let table = match self
76                .config
77                .catalog_manager
78                .table(&name[0], &name[1], &name[2], None)
79                .await
80            {
81                Ok(Some(table)) => table,
82                Ok(None) => {
83                    return Err(UnexpectedSnafu {
84                        reason: format!(
85                            "Flow {} source table {} not found for sequence_range capability check",
86                            self.config.flow_id,
87                            name.join(".")
88                        ),
89                    }
90                    .build());
91                }
92                Err(err) => Err(BoxedError::new(err)).context(ExternalSnafu)?,
93            };
94
95            let info = table.table_info();
96            let preserves = info.meta.engine == "mito"
97                && info
98                    .meta
99                    .options
100                    .extra_options
101                    .get(PRESERVE_ROW_SEQUENCE)
102                    .is_some_and(|value| value.eq_ignore_ascii_case("true"));
103            if !preserves {
104                return Ok(false);
105            }
106        }
107        Ok(!self.config.source_table_names.is_empty())
108    }
109
110    /// For incremental-mode SQL queries, attempt to prepare an executable plan
111    /// that is safe for incremental scan extensions.
112    ///
113    /// Returns `Some(plan)` when incremental extensions are safe. For an
114    /// incremental-delta query, `None` means the caller must not execute the
115    /// original plan without incremental extensions, because that would become
116    /// an unfiltered full snapshot; the caller should restore the dirty signal
117    /// and skip the current round instead. The returned plan may be either a
118    /// rewritten delta-LEFT-JOIN-sink merge plan or the original plan. In
119    /// particular, plain GROUP BY queries with no aggregate merge columns are
120    /// incremental safe without a rewrite, so they return `Some(original_plan)`.
121    pub(super) async fn prepare_plan_for_incremental(
122        &self,
123        engine: &QueryEngineRef,
124        plan: &LogicalPlan,
125    ) -> Result<Option<LogicalPlan>, Error> {
126        let is_incremental_sql = {
127            let state = self.state.read().unwrap();
128            if state.is_incremental_disabled() {
129                return Ok(None);
130            }
131            state.checkpoint_mode() == CheckpointMode::Incremental
132                && matches!(self.config.query_type, QueryType::Sql)
133        };
134
135        if !is_incremental_sql {
136            return Ok(None);
137        }
138
139        // Extract inner query plan from the DML wrapper.
140        // Non-DML or non-SQL plans bypass the rewrite and keep checkpoint mode;
141        // non-aggregate TQL or non-INSERT plans do not need incremental scan extensions.
142        let inner_plan = match plan {
143            LogicalPlan::Dml(dml) => dml.input.as_ref().clone(),
144            _ => return Ok(None),
145        };
146
147        // Analyze the plan for incremental rewritability.
148        // Incremental reads currently require aggregate / group-by plans that
149        // can be rewritten into a delta-left-join-sink merge. Non-aggregate SQL
150        // (projection, filter, or other non-aggregate shapes) stays full-snapshot
151        // until separately supported, and incremental mode is permanently
152        // disabled for this flow.
153        let Some(analysis) = analyze_incremental_aggregate_plan(&inner_plan)? else {
154            warn!(
155                "Flow {} incremental mode but plan is not an aggregate query; \
156                 permanently disabling incremental for this flow",
157                self.config.flow_id
158            );
159            if self.config.exact_sequence_range_required {
160                return Err(UnexpectedSnafu {
161                    reason: format!(
162                        "Flow {} requires exact sequence-range reads, but its incremental plan is unsupported",
163                        self.config.flow_id
164                    ),
165                }
166                .build());
167            }
168            self.state.write().unwrap().disable_incremental();
169            return Ok(None);
170        };
171
172        if !analysis.unsupported_exprs.is_empty() {
173            warn!(
174                "Flow {} incremental aggregate contains unsupported expressions {:?}; \
175                 permanently disabling incremental for this flow",
176                self.config.flow_id, analysis.unsupported_exprs
177            );
178            if self.config.exact_sequence_range_required {
179                return Err(UnexpectedSnafu {
180                    reason: format!(
181                        "Flow {} requires exact sequence-range reads, but its incremental aggregate is unsupported: {:?}",
182                        self.config.flow_id, analysis.unsupported_exprs
183                    ),
184                }
185                .build());
186            }
187            self.state.write().unwrap().disable_incremental();
188            return Ok(None);
189        }
190
191        // Plain GROUP BY without aggregate expressions has no values to
192        // merge between delta and sink. The incremental delta scan emits
193        // changed groups, and sink primary-key write semantics make this
194        // idempotent; no explicit left-join rewrite is needed.
195        if analysis.merge_columns.is_empty() {
196            return Ok(Some(plan.clone()));
197        }
198
199        // Fetch sink table for the merge rewrite.
200        // Transient errors (catalog, schema, filter, or rewrite) should not
201        // permanently disable incremental mode. They also must not execute the
202        // original plan without incremental extensions, because that would be an
203        // unfiltered full snapshot. The caller will restore the dirty signal and
204        // skip this round while keeping incremental retryable.
205        let sink_table = match get_table_info_df_schema(
206            self.config.catalog_manager.clone(),
207            self.config.sink_table_name.clone(),
208        )
209        .await
210        {
211            Ok((table, _)) => table,
212            Err(err) => {
213                warn!(
214                    "Flow {} failed to fetch sink table for incremental rewrite; \
215                     skipping this round to avoid unfiltered full snapshot: {:?}",
216                    self.config.flow_id, err
217                );
218                return Ok(None);
219            }
220        };
221        let rewritten_inner = match rewrite_incremental_aggregate_with_sink_merge(
222            &inner_plan,
223            &analysis,
224            engine,
225            sink_table,
226            &self.config.sink_table_name,
227            None,
228        )
229        .await
230        {
231            Ok(plan) => plan,
232            Err(err) => {
233                warn!(
234                    "Flow {} failed to rewrite incremental aggregate with sink merge; \
235                     skipping this round to avoid unfiltered full snapshot: {:?}",
236                    self.config.flow_id, err
237                );
238                return Ok(None);
239            }
240        };
241
242        // Reconstruct DML plan with the rewritten inner plan
243        let rewritten = match plan {
244            LogicalPlan::Dml(dml) => LogicalPlan::Dml(DmlStatement::new(
245                dml.table_name.clone(),
246                dml.target.clone(),
247                dml.op.clone(),
248                Arc::new(rewritten_inner),
249            )),
250            _ => unreachable!("already matched Dml above"),
251        };
252
253        debug!(
254            "Flow {} rewrote incremental SQL aggregate query with sink merge",
255            self.config.flow_id
256        );
257
258        Ok(Some(rewritten))
259    }
260
261    pub(super) async fn build_flow_query_extensions(
262        &self,
263        incremental_safe: bool,
264        can_advance_checkpoints: bool,
265    ) -> Result<Vec<(&'static str, String)>, Error> {
266        let mut extensions = vec![("flow.return_region_seq", "true".to_string())];
267
268        let incremental_checkpoints_json = {
269            let state = self.state.read().unwrap();
270            if incremental_safe
271                && can_advance_checkpoints
272                && !state.is_incremental_disabled()
273                && state.checkpoint_mode() == CheckpointMode::Incremental
274                && !state.checkpoints().is_empty()
275            {
276                Some(serde_json::to_string(state.checkpoints()).map_err(|err| {
277                    UnexpectedSnafu {
278                        reason: format!("Failed to serialize checkpoint map: {err}"),
279                    }
280                    .build()
281                })?)
282            } else {
283                None
284            }
285        };
286
287        if let Some(checkpoints_json) = incremental_checkpoints_json {
288            // Select `sequence_range` only when every append-only source table
289            // proves the `preserve_row_sequence` capability; otherwise retain
290            // the historical `memtable_only` mode. The `sequence_range` scan
291            // keeps SSTs and reads the exact (checkpoint, scan-open snapshot]
292            // row-level delta; the engine fails closed when the capability
293            // does not hold at scan time.
294            let capable = match self.sequence_range_capable().await {
295                Ok(capable) => capable,
296                Err(err) => {
297                    if self.config.exact_sequence_range_required {
298                        return Err(err);
299                    }
300                    false
301                }
302            };
303            if self.config.exact_sequence_range_required && !capable {
304                return Err(UnexpectedSnafu {
305                    reason: format!(
306                        "Flow {} requires exact sequence-range reads, but source capability was revoked",
307                        self.config.flow_id
308                    ),
309                }
310                .build());
311            }
312            let sink_table_id = self.sink_table_id().await?;
313            let incremental_mode = if capable {
314                debug!(
315                    "Flow {} selected sequence_range incremental mode",
316                    self.config.flow_id
317                );
318                FLOW_INCREMENTAL_MODE_SEQUENCE_RANGE
319            } else {
320                FLOW_INCREMENTAL_MODE_MEMTABLE_ONLY
321            };
322            extensions.push((FLOW_SINK_TABLE_ID, sink_table_id.to_string()));
323            extensions.push((FLOW_INCREMENTAL_MODE, incremental_mode.to_string()));
324            extensions.push((FLOW_INCREMENTAL_AFTER_SEQS, checkpoints_json));
325        }
326
327        Ok(extensions)
328    }
329}