flow/batching_mode/task/
inc.rs1use 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 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 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 let inner_plan = match plan {
143 LogicalPlan::Dml(dml) => dml.input.as_ref().clone(),
144 _ => return Ok(None),
145 };
146
147 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 if analysis.merge_columns.is_empty() {
196 return Ok(Some(plan.clone()));
197 }
198
199 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 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 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}