1use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
16use std::sync::{Arc, RwLock};
17use std::time::{Duration, SystemTime, UNIX_EPOCH};
18
19use api::v1::{CreateTableExpr, TableName};
20use catalog::CatalogManagerRef;
21use catalog::kvbackend::KvBackendCatalogManager;
22use client::OutputWithMetrics;
23use common_error::ext::BoxedError;
24use common_meta::key::schema_name::SchemaNameKey;
25use common_query::OutputData;
26use common_query::logical_plan::breakup_insert_plan;
27use common_telemetry::tracing::warn;
28use common_telemetry::{debug, info};
29use common_time::{TimeToLive, Timestamp};
30use datafusion::datasource::DefaultTableSource;
31use datafusion::sql::unparser::expr_to_sql;
32use datafusion_common::tree_node::{Transformed, TreeNode};
33use datafusion_common::utils::quote_identifier;
34use datafusion_common::{DFSchemaRef, ScalarValue, TableReference};
35use datafusion_expr::{DmlStatement, LogicalPlan, Projection, WriteOp, col, lit};
36use datatypes::schema::Schema;
37use datatypes::vectors::Helper;
38use futures::TryStreamExt;
39use query::QueryEngineRef;
40use query::options::{
41 FLOW_INCREMENTAL_AFTER_SEQS, FLOW_INCREMENTAL_MODE, FLOW_INCREMENTAL_MODE_SEQUENCE_RANGE,
42 FLOW_RETURN_REGION_SEQ,
43};
44use query::query_engine::DefaultSerializer;
45use session::context::QueryContextRef;
46use snafu::{OptionExt, ResultExt};
47use sql::parsers::utils::is_tql;
48use store_api::mito_engine_options::MERGE_MODE_KEY;
49use substrait::{DFLogicalSubstraitConvertor, SubstraitPlan};
50use table::TableRef;
51use table::table::adapter::DfTableProviderAdapter;
52use tokio::sync::oneshot::error::TryRecvError;
53use tokio::sync::{Mutex, OwnedMutexGuard, oneshot};
54use tokio::time::Instant;
55
56use crate::batching_mode::BatchingModeOptions;
57use crate::batching_mode::checkpoint::checkpoint_mode_label;
58use crate::batching_mode::eval_schedule::{EvalSchedule, select_due_scheduled_times};
59use crate::batching_mode::frontend_client::{FrontendClient, PeerDesc};
60use crate::batching_mode::state::{
61 CheckpointMode, DirtyTimeWindows, FilterExprInfo, TaskState, to_df_literal,
62};
63use crate::batching_mode::table_creator::{QueryType, create_table_with_expr};
64use crate::batching_mode::time_window::TimeWindowExpr;
65use crate::batching_mode::utils::{
66 AddFilterRewriter, ColumnMatcherRewriter, analyze_incremental_aggregate_plan, df_plan_to_sql,
67 gen_plan_with_matching_schema_and_values, get_table_info_df_schema, sql_to_df_plan,
68};
69use crate::df_optimizer::apply_df_optimizer;
70use crate::error::{
71 DatafusionSnafu, ExternalSnafu, InvalidQuerySnafu, SubstraitEncodeLogicalPlanSnafu,
72 UnexpectedSnafu,
73};
74use crate::metrics::{
75 METRIC_FLOW_BATCHING_ENGINE_ERROR_CNT, METRIC_FLOW_BATCHING_ENGINE_QUERY_TIME,
76 METRIC_FLOW_BATCHING_ENGINE_SLOW_QUERY, METRIC_FLOW_BATCHING_ENGINE_START_QUERY_CNT,
77 METRIC_FLOW_ROWS,
78};
79use crate::{Error, FlowId};
80
81mod ckpt;
82mod inc;
83
84fn wall_clock_unix_secs() -> i64 {
86 SystemTime::now()
87 .duration_since(UNIX_EPOCH)
88 .unwrap_or_default()
89 .as_secs() as i64
90}
91
92fn initial_schedule_cursor(start_secs: i64, interval_secs: i64) -> Result<i64, Error> {
100 let cursor = i128::from(start_secs) - i128::from(interval_secs);
101 i64::try_from(cursor).map_err(|_| {
102 UnexpectedSnafu {
103 reason: format!(
104 "Cannot compute the initial eval schedule cursor one interval before start {start_secs} (interval={interval_secs}): {cursor} does not fit in i64"
105 ),
106 }
107 .build()
108 })
109}
110
111fn sleep_delta_secs(next: i64, wall_now_secs: i64) -> Result<u64, Error> {
119 let delta = i128::from(next) - i128::from(wall_now_secs);
120 u64::try_from(delta).map_err(|_| {
121 UnexpectedSnafu {
122 reason: format!(
123 "Cannot sleep until the next scheduled time {next}: the delta from wall clock {wall_now_secs} is {delta} seconds, which does not fit in u64"
124 ),
125 }
126 .build()
127 })
128}
129
130fn scheduled_time_millis(scheduled_time_secs: i64) -> Result<i64, Error> {
137 scheduled_time_secs.checked_mul(1000).ok_or_else(|| {
138 UnexpectedSnafu {
139 reason: format!(
140 "Cannot convert scheduled time {scheduled_time_secs}s to milliseconds: the product exceeds i64"
141 ),
142 }
143 .build()
144 })
145}
146
147#[derive(Clone)]
149pub struct TaskConfig {
150 pub flow_id: FlowId,
151 pub query: String,
152 pub output_schema: DFSchemaRef,
154 pub time_window_expr: Option<TimeWindowExpr>,
155 pub expire_after: Option<i64>,
157 pub sink_table_name: [String; 3],
158 pub source_table_names: HashSet<[String; 3]>,
159 pub catalog_manager: CatalogManagerRef,
160 pub query_type: QueryType,
161 pub batch_opts: Arc<BatchingModeOptions>,
162 pub exact_sequence_range_required: bool,
163 pub flow_eval_interval: Option<Duration>,
164 pub eval_schedule: Option<EvalSchedule>,
166}
167
168fn determine_query_type(query: &str, query_ctx: &QueryContextRef) -> Result<QueryType, Error> {
169 let is_tql = is_tql(query_ctx.sql_dialect(), query)
170 .map_err(BoxedError::new)
171 .context(ExternalSnafu)?;
172 Ok(if is_tql {
173 QueryType::Tql
174 } else {
175 QueryType::Sql
176 })
177}
178
179fn is_merge_mode_last_non_null(options: &HashMap<String, String>) -> bool {
180 options
181 .get(MERGE_MODE_KEY)
182 .map(|mode| mode.eq_ignore_ascii_case("last_non_null"))
183 .unwrap_or(false)
184}
185
186fn encode_insert_plan_request(
187 insert_to: TableName,
188 insert_input_plan: &LogicalPlan,
189) -> Result<api::v1::QueryRequest, Error> {
190 let message = DFLogicalSubstraitConvertor {}
191 .encode(insert_input_plan, DefaultSerializer)
192 .context(SubstraitEncodeLogicalPlanSnafu)?;
193 Ok(api::v1::QueryRequest {
194 query: Some(api::v1::query_request::Query::InsertIntoPlan(
195 api::v1::InsertIntoPlan {
196 table_name: Some(insert_to),
197 logical_plan: message.to_vec(),
198 },
199 )),
200 })
201}
202
203fn recovery_aggregate_input(plan: &LogicalPlan) -> Result<LogicalPlan, Error> {
204 let plan = match plan {
205 LogicalPlan::Projection(projection) => projection.input.as_ref(),
206 _ => plan,
207 };
208 let LogicalPlan::Aggregate(aggregate) = plan else {
209 return UnexpectedSnafu {
210 reason: "Recovery timestamp projection did not find an aggregate".to_string(),
211 }
212 .fail();
213 };
214 Ok(aggregate.input.as_ref().clone())
215}
216
217fn capture_recovery_batch_windows(
218 batch: &common_recordbatch::RecordBatch,
219 time_window_expr: &TimeWindowExpr,
220 windows: &mut BTreeSet<(Timestamp, Timestamp)>,
221) -> Result<(), Error> {
222 if batch.num_columns() != 1 {
223 return UnexpectedSnafu {
224 reason: format!(
225 "Recovery timestamp projection returned {} columns instead of one",
226 batch.num_columns()
227 ),
228 }
229 .fail();
230 }
231 let values = Helper::try_into_vector(batch.column(0).clone())
232 .map_err(BoxedError::new)
233 .context(ExternalSnafu)?;
234 for index in 0..values.len() {
235 let timestamp = values.get(index).as_timestamp().context(UnexpectedSnafu {
236 reason: "Recovery timestamp projection returned a null or non-timestamp value"
237 .to_string(),
238 })?;
239 let (start, end) = time_window_expr.eval(timestamp)?;
240 let window = (
241 start.context(UnexpectedSnafu {
242 reason: "Recovery time-window expression returned no lower bound".to_string(),
243 })?,
244 end.context(UnexpectedSnafu {
245 reason: "Recovery time-window expression returned no upper bound".to_string(),
246 })?,
247 );
248 windows.insert(window);
249 }
250 Ok(())
251}
252
253fn format_insert_target_columns(plan: &LogicalPlan) -> String {
254 plan.schema()
255 .fields()
256 .iter()
257 .map(|field| quote_identifier(field.name()).to_string())
258 .collect::<Vec<_>>()
259 .join(", ")
260}
261
262pub struct BatchingExecutionGuard {
265 _lock: OwnedMutexGuard<()>,
266 restore: Option<(Arc<RwLock<TaskState>>, QueryContextRef)>,
267}
268
269impl BatchingExecutionGuard {
270 fn new(lock: OwnedMutexGuard<()>) -> Self {
271 Self {
272 _lock: lock,
273 restore: None,
274 }
275 }
276
277 fn restore_query_context(&mut self, state: Arc<RwLock<TaskState>>, old_ctx: QueryContextRef) {
278 self.restore = Some((state, old_ctx));
279 }
280}
281
282impl Drop for BatchingExecutionGuard {
283 fn drop(&mut self) {
284 if let Some((state, old_ctx)) = self.restore.take()
285 && let Ok(mut state) = state.write()
286 {
287 state.query_ctx = old_ctx;
288 }
289 }
290}
291
292#[derive(Clone)]
293pub struct BatchingTask {
294 pub config: Arc<TaskConfig>,
295 pub state: Arc<RwLock<TaskState>>,
296 execution_lock: Arc<Mutex<()>>,
300 execution: Option<Arc<dyn crate::BatchingExecution>>,
301}
302
303pub struct TaskArgs<'a> {
305 pub flow_id: FlowId,
306 pub query: &'a str,
307 pub plan: LogicalPlan,
308 pub time_window_expr: Option<TimeWindowExpr>,
309 pub expire_after: Option<i64>,
310 pub sink_table_name: [String; 3],
311 pub source_table_names: Vec<[String; 3]>,
312 pub query_ctx: QueryContextRef,
313 pub catalog_manager: CatalogManagerRef,
314 pub shutdown_rx: oneshot::Receiver<()>,
315 pub batch_opts: Arc<BatchingModeOptions>,
316 pub flow_eval_interval: Option<Duration>,
317 pub eval_schedule: Option<EvalSchedule>,
319}
320
321pub struct PlanInfo {
322 pub plan: LogicalPlan,
323 pub dirty_restore: DirtyRestore,
324 pub coverage: QueryCoverage,
325}
326
327#[derive(Clone)]
328pub enum QueryCoverage {
329 UnfilteredFull,
334 ScopedBaseRepair,
337 FencedRepairChunk { high: BTreeMap<u64, u64> },
341 IncrementalDelta,
343}
344
345impl QueryCoverage {
346 fn is_incremental_delta(&self) -> bool {
349 matches!(self, Self::IncrementalDelta)
350 }
351
352 fn snapshot_seqs(&self) -> HashMap<u64, u64> {
355 match self {
356 Self::FencedRepairChunk { high } => high.iter().map(|(k, v)| (*k, *v)).collect(),
357 _ => HashMap::new(),
358 }
359 }
360}
361
362pub enum DirtyRestore {
363 Scoped(FilterExprInfo),
366 Unscoped(DirtyTimeWindows),
373}
374
375pub struct ExecuteOnceOutcome {
376 pub new_query: Option<PlanInfo>,
377 pub result: Result<Option<(usize, Duration)>, Error>,
383}
384
385impl BatchingTask {
386 #[allow(clippy::too_many_arguments)]
387 pub fn try_new(args: TaskArgs<'_>) -> Result<Self, Error> {
388 Self::try_new_with_exact_sequence_range_required(args, false)
389 }
390
391 pub fn try_new_with_exact_sequence_range_required(
392 TaskArgs {
393 flow_id,
394 query,
395 plan,
396 time_window_expr,
397 expire_after,
398 sink_table_name,
399 source_table_names,
400 query_ctx,
401 catalog_manager,
402 shutdown_rx,
403 batch_opts,
404 flow_eval_interval,
405 eval_schedule,
406 }: TaskArgs<'_>,
407 exact_sequence_range_required: bool,
408 ) -> Result<Self, Error> {
409 let mut state = TaskState::with_dirty_time_windows(
410 query_ctx.clone(),
411 shutdown_rx,
412 DirtyTimeWindows::new(
413 batch_opts.experimental_max_filter_num_per_query,
414 batch_opts.experimental_time_window_merge_threshold,
415 ),
416 );
417 if !batch_opts.experimental_enable_incremental_read {
418 state.disable_incremental();
419 }
420
421 Ok(Self {
422 config: Arc::new(TaskConfig {
423 flow_id,
424 query: query.to_string(),
425 time_window_expr,
426 expire_after,
427 sink_table_name,
428 source_table_names: source_table_names.into_iter().collect(),
429 catalog_manager,
430 output_schema: plan.schema().clone(),
431 query_type: determine_query_type(query, &query_ctx)?,
432 exact_sequence_range_required,
433 batch_opts,
434 flow_eval_interval,
435 eval_schedule,
436 }),
437 state: Arc::new(RwLock::new(state)),
438 execution_lock: Arc::new(Mutex::new(())),
439 execution: None,
440 })
441 }
442
443 pub(crate) fn with_execution(
444 mut self,
445 execution: Option<Arc<dyn crate::BatchingExecution>>,
446 ) -> Self {
447 self.execution = execution;
448 self
449 }
450
451 pub(crate) fn stop_execution(&self) {
452 if let Some(execution) = &self.execution {
453 execution.stop();
454 }
455 }
456
457 pub fn last_execution_time_millis(&self) -> Option<i64> {
458 self.state.read().unwrap().last_execution_time_millis()
459 }
460
461 pub fn start_time_millis(&self) -> Option<i64> {
462 self.state.read().unwrap().start_time_millis()
463 }
464
465 fn frontend_extensions(&self) -> HashMap<String, String> {
468 let ctx = self.state.read().unwrap();
469 let all = ctx.query_ctx.extensions();
470 let mut flow_exts = HashMap::new();
471 if let Some(v) = all.get(query::options::FLOW_SCHEDULED_TIME_MILLIS) {
474 flow_exts.insert(
475 query::options::FLOW_SCHEDULED_TIME_MILLIS.to_string(),
476 v.clone(),
477 );
478 }
479 flow_exts
480 }
481
482 pub fn mark_all_windows_as_dirty(&self) -> Result<(), Error> {
486 let now = SystemTime::now();
487 let now = Timestamp::new_second(
488 now.duration_since(UNIX_EPOCH)
489 .expect("Time went backwards")
490 .as_secs() as _,
491 );
492 let lower_bound = self
493 .config
494 .expire_after
495 .map(|e| now.sub_duration(Duration::from_secs(e as _)))
496 .transpose()
497 .map_err(BoxedError::new)
498 .context(ExternalSnafu)?
499 .unwrap_or(Timestamp::new_second(0));
500 debug!(
501 "Flow {} mark range ({:?}, {:?}) as dirty",
502 self.config.flow_id, lower_bound, now
503 );
504 self.state
505 .write()
506 .unwrap()
507 .dirty_time_windows
508 .add_window(lower_bound, Some(now));
509 Ok(())
510 }
511
512 pub async fn check_or_create_sink_table(
514 &self,
515 engine: &QueryEngineRef,
516 frontend_client: &Arc<FrontendClient>,
517 ) -> Result<TableRef, Error> {
518 if !self.is_table_exist(&self.config.sink_table_name).await? {
519 let create_table = self.gen_create_table_expr(engine.clone()).await?;
520 info!(
521 "Try creating sink table(if not exists) with expr: {:?}",
522 create_table
523 );
524 self.create_table(frontend_client, create_table).await?;
525 info!(
526 "Sink table {}(if not exists) created",
527 self.config.sink_table_name.join(".")
528 );
529 }
530 let (table, _) = get_table_info_df_schema(
531 self.config.catalog_manager.clone(),
532 self.config.sink_table_name.clone(),
533 )
534 .await?;
535 Ok(table)
536 }
537
538 pub fn query_context_snapshot(&self) -> QueryContextRef {
540 let query_ctx = self.state.read().unwrap().query_ctx.clone();
541 Arc::new(query_ctx.fork())
542 }
543
544 pub fn recovery_retention_lower_bound(&self) -> Result<Option<Timestamp>, Error> {
546 let Some(expire_after) = self.config.expire_after else {
547 return Ok(None);
548 };
549 let expire_after = u64::try_from(expire_after).map_err(|_| {
550 UnexpectedSnafu {
551 reason: format!(
552 "Flow {} has negative expire_after {expire_after}",
553 self.config.flow_id
554 ),
555 }
556 .build()
557 })?;
558 let now = Timestamp::new_second(
559 SystemTime::now()
560 .duration_since(UNIX_EPOCH)
561 .map_err(|err| {
562 UnexpectedSnafu {
563 reason: format!("Failed to read recovery wall clock: {err}"),
564 }
565 .build()
566 })?
567 .as_secs() as i64,
568 );
569 let lower = now
570 .sub_duration(Duration::from_secs(expire_after))
571 .map_err(BoxedError::new)
572 .context(ExternalSnafu)?;
573 self.config
574 .time_window_expr
575 .as_ref()
576 .context(UnexpectedSnafu {
577 reason: "Recovery expiry requires a time-window expression".to_string(),
578 })?
579 .eval(lower)?
580 .0
581 .context(UnexpectedSnafu {
582 reason: "Recovery expiry time-window expression returned no lower bound"
583 .to_string(),
584 })
585 .map(Some)
586 }
587
588 pub async fn validate_recovery_retention(
590 &self,
591 retention_lower: Option<Timestamp>,
592 windows: &[(Timestamp, Timestamp)],
593 ) -> Result<(), Error> {
594 for name in &self.config.source_table_names {
595 let table = self
596 .config
597 .catalog_manager
598 .table(&name[0], &name[1], &name[2], None)
599 .await
600 .map_err(BoxedError::new)
601 .context(ExternalSnafu)?
602 .context(UnexpectedSnafu {
603 reason: format!(
604 "Flow {} source table {} is unavailable for recovery retention validation",
605 self.config.flow_id,
606 name.join(".")
607 ),
608 })?;
609 let ttl = if let Some(ttl) = table.table_info().meta.options.ttl {
610 ttl
611 } else {
612 let manager = self
613 .config
614 .catalog_manager
615 .as_any()
616 .downcast_ref::<KvBackendCatalogManager>()
617 .context(UnexpectedSnafu {
618 reason: format!(
619 "Flow {} cannot resolve inherited TTL for source table {} during recovery",
620 self.config.flow_id,
621 name.join(".")
622 ),
623 })?;
624 manager
625 .table_metadata_manager_ref()
626 .schema_manager()
627 .get(SchemaNameKey::new(&name[0], &name[1]))
628 .await
629 .map_err(BoxedError::new)
630 .context(ExternalSnafu)?
631 .context(UnexpectedSnafu {
632 reason: format!(
633 "Flow {} schema {}.{} is unavailable for recovery retention validation",
634 self.config.flow_id, name[0], name[1]
635 ),
636 })?
637 .ttl
638 .map(Into::into)
639 .unwrap_or(TimeToLive::Forever)
640 };
641
642 match ttl {
643 TimeToLive::Forever => {}
644 TimeToLive::Instant => {
645 return UnexpectedSnafu {
646 reason: format!(
647 "Flow {} source table {} has instant TTL and cannot be recovered",
648 self.config.flow_id,
649 name.join(".")
650 ),
651 }
652 .fail();
653 }
654 TimeToLive::Duration(ttl) => {
655 let expire_after = self.config.expire_after.context(UnexpectedSnafu {
656 reason: format!(
657 "Flow {} source table {} has TTL {ttl:?}, but recovery has no expire_after",
658 self.config.flow_id,
659 name.join(".")
660 ),
661 })?;
662 let expire_after = u64::try_from(expire_after).map_err(|_| {
663 UnexpectedSnafu {
664 reason: format!(
665 "Flow {} has negative expire_after {expire_after} during recovery retention validation",
666 self.config.flow_id
667 ),
668 }
669 .build()
670 })?;
671 if ttl <= Duration::from_secs(expire_after) {
672 return UnexpectedSnafu {
673 reason: format!(
674 "Flow {} source table {} TTL {ttl:?} must exceed expire_after {expire_after}s for recovery",
675 self.config.flow_id,
676 name.join(".")
677 ),
678 }
679 .fail();
680 }
681 let mut oldest = retention_lower.context(UnexpectedSnafu {
682 reason: format!(
683 "Flow {} source table {} has finite TTL {ttl:?}, but recovery has no retained lower bound",
684 self.config.flow_id,
685 name.join(".")
686 ),
687 })?;
688 for (start, _) in windows {
689 oldest = oldest.min(*start);
690 }
691 let cutoff = Timestamp::current_millis()
692 .sub_duration(ttl)
693 .map_err(BoxedError::new)
694 .context(ExternalSnafu)?;
695 if oldest <= cutoff {
696 return UnexpectedSnafu {
697 reason: format!(
698 "Flow {} recovery requires source table {} data from {oldest:?}, older than TTL {ttl:?} cutoff {cutoff:?}",
699 self.config.flow_id,
700 name.join(".")
701 ),
702 }
703 .fail();
704 }
705 }
706 }
707 }
708 Ok(())
709 }
710
711 pub async fn capture_recovery_windows(
715 &self,
716 engine: &QueryEngineRef,
717 frontend_client: &FrontendClient,
718 lower: &BTreeMap<u64, u64>,
719 ) -> Result<(BTreeMap<u64, u64>, Vec<(Timestamp, Timestamp)>), Error> {
720 let retention_lower = self.recovery_retention_lower_bound()?;
721 self.capture_recovery_windows_since(engine, frontend_client, lower, retention_lower)
722 .await
723 }
724
725 pub async fn capture_recovery_windows_since(
731 &self,
732 engine: &QueryEngineRef,
733 frontend_client: &FrontendClient,
734 lower: &BTreeMap<u64, u64>,
735 retention_lower: Option<Timestamp>,
736 ) -> Result<(BTreeMap<u64, u64>, Vec<(Timestamp, Timestamp)>), Error> {
737 if lower.is_empty() {
738 return UnexpectedSnafu {
739 reason: format!(
740 "Flow {} recovery window capture requires nonempty lower sequence bounds",
741 self.config.flow_id
742 ),
743 }
744 .fail();
745 }
746 if !self.sequence_range_capable().await? {
747 return UnexpectedSnafu {
748 reason: format!(
749 "Flow {} recovery window capture requires sequence_range-capable sources",
750 self.config.flow_id
751 ),
752 }
753 .fail();
754 }
755
756 let query_ctx = self.query_context_snapshot();
757 let plan = sql_to_df_plan(query_ctx, engine.clone(), &self.config.query, false).await?;
758 let input = recovery_aggregate_input(&plan)?;
759 let Some(analysis) = analyze_incremental_aggregate_plan(&plan)? else {
760 return UnexpectedSnafu {
761 reason: format!(
762 "Flow {} recovery window capture requires a supported aggregate input",
763 self.config.flow_id
764 ),
765 }
766 .fail();
767 };
768 if !analysis.unsupported_exprs.is_empty() {
769 return UnexpectedSnafu {
770 reason: format!(
771 "Flow {} recovery window capture has unsupported aggregate expressions: {:?}",
772 self.config.flow_id, analysis.unsupported_exprs
773 ),
774 }
775 .fail();
776 }
777 let time_window_expr = self
778 .config
779 .time_window_expr
780 .as_ref()
781 .context(UnexpectedSnafu {
782 reason: format!(
783 "Flow {} recovery window capture requires a time-window expression",
784 self.config.flow_id
785 ),
786 })?;
787 let input = if let Some(retention_lower) = retention_lower {
788 let mut add_filter = AddFilterRewriter::new(
789 col(&time_window_expr.column_name).gt_eq(lit(to_df_literal(retention_lower)?)),
790 );
791 input
792 .rewrite(&mut add_filter)
793 .with_context(|_| DatafusionSnafu {
794 context: "Failed to apply recovery expire_after filter".to_string(),
795 })?
796 .data
797 } else {
798 input
799 };
800 let timestamp_plan = LogicalPlan::Projection(
801 Projection::try_new(vec![col(&time_window_expr.column_name)], Arc::new(input))
802 .context(DatafusionSnafu {
803 context: "Failed to project recovery source timestamps".to_string(),
804 })?,
805 );
806 let catalog = &self.config.sink_table_name[0];
807 let schema = &self.config.sink_table_name[1];
808 let timestamp_plan = timestamp_plan
809 .clone()
810 .transform_down_with_subqueries(|p| {
811 if let LogicalPlan::TableScan(mut table_scan) = p {
812 let resolved = table_scan.table_name.resolve(catalog, schema);
813 table_scan.table_name = resolved.into();
814 Ok(Transformed::yes(LogicalPlan::TableScan(table_scan)))
815 } else {
816 Ok(Transformed::no(p))
817 }
818 })
819 .with_context(|_| DatafusionSnafu {
820 context: format!(
821 "Failed to fix table ref in recovery timestamp plan, plan={:?}",
822 timestamp_plan
823 ),
824 })?
825 .data;
826 let message = DFLogicalSubstraitConvertor {}
827 .encode(×tamp_plan, DefaultSerializer)
828 .context(SubstraitEncodeLogicalPlanSnafu)?;
829 let request = api::v1::QueryRequest {
830 query: Some(api::v1::query_request::Query::LogicalPlan(message.to_vec())),
831 };
832 let lower_json = serde_json::to_string(lower).map_err(|err| {
833 UnexpectedSnafu {
834 reason: format!("Failed to serialize recovery lower sequence bounds: {err}"),
835 }
836 .build()
837 })?;
838 let extensions = [
839 (FLOW_RETURN_REGION_SEQ, "true"),
840 (FLOW_INCREMENTAL_MODE, FLOW_INCREMENTAL_MODE_SEQUENCE_RANGE),
841 (FLOW_INCREMENTAL_AFTER_SEQS, lower_json.as_str()),
842 ];
843 let mut peer_desc = None;
844 let result = frontend_client
845 .query_with_terminal_metrics(
846 catalog,
847 schema,
848 request,
849 &extensions,
850 &HashMap::new(),
851 &mut peer_desc,
852 )
853 .await?;
854 let mut aligned_windows = BTreeSet::new();
855 let metrics = result.metrics.clone();
856 match result.output.data {
857 OutputData::AffectedRows(_) => {
858 return UnexpectedSnafu {
859 reason: "Recovery timestamp projection unexpectedly returned affected rows"
860 .to_string(),
861 }
862 .fail();
863 }
864 OutputData::RecordBatches(batches) => {
865 for batch in batches.iter() {
866 capture_recovery_batch_windows(batch, time_window_expr, &mut aligned_windows)?;
867 }
868 }
869 OutputData::Stream(mut stream) => {
870 while let Some(batch) = stream
871 .try_next()
872 .await
873 .map_err(BoxedError::new)
874 .context(ExternalSnafu)?
875 {
876 capture_recovery_batch_windows(&batch, time_window_expr, &mut aligned_windows)?;
877 }
878 }
879 }
880 if !metrics.is_ready() {
881 return UnexpectedSnafu {
882 reason: "Recovery timestamp projection ended without terminal metrics".to_string(),
883 }
884 .fail();
885 }
886 let participating = metrics.participating_regions().context(UnexpectedSnafu {
887 reason: "Recovery timestamp projection has no participating-region proof".to_string(),
888 })?;
889 let high = metrics.region_watermark_map().context(UnexpectedSnafu {
890 reason: "Recovery timestamp projection has no terminal watermark proof".to_string(),
891 })?;
892 if participating.len() != lower.len()
893 || high.len() != lower.len()
894 || participating
895 .iter()
896 .any(|region| !lower.contains_key(region))
897 || lower
898 .iter()
899 .any(|(region, low)| high.get(region).is_none_or(|watermark| watermark < low))
900 || high.keys().any(|region| !participating.contains(region))
901 {
902 return UnexpectedSnafu {
903 reason: format!(
904 "Recovery timestamp projection returned incomplete or regressing terminal proof: lower={lower:?}, participating={participating:?}, high={high:?}"
905 ),
906 }
907 .fail();
908 }
909 Ok((
910 high.into_iter().collect(),
911 aligned_windows.into_iter().collect(),
912 ))
913 }
914
915 pub async fn validate_sink_table_schema(&self, engine: &QueryEngineRef) -> Result<(), Error> {
917 self.validate_sink_table_schema_with_values(engine, &BTreeMap::new())
918 .await
919 }
920
921 pub async fn validate_sink_table_schema_with_values(
923 &self,
924 engine: &QueryEngineRef,
925 values: &BTreeMap<String, ScalarValue>,
926 ) -> Result<(), Error> {
927 let (table, _) = get_table_info_df_schema(
928 self.config.catalog_manager.clone(),
929 self.config.sink_table_name.clone(),
930 )
931 .await?;
932
933 let table_meta = &table.table_info().meta;
934 let merge_mode_last_non_null =
935 is_merge_mode_last_non_null(&table_meta.options.extra_options);
936 let primary_key_indices = table_meta.primary_key_indices.clone();
937
938 gen_plan_with_matching_schema_and_values(
939 &self.config.query,
940 self.query_context_snapshot(),
941 engine.clone(),
942 table_meta.schema.clone(),
943 &primary_key_indices,
944 merge_mode_last_non_null,
945 values,
946 )
947 .await
948 .map(|_| ())
949 }
950
951 async fn is_table_exist(&self, table_name: &[String; 3]) -> Result<bool, Error> {
952 self.config
953 .catalog_manager
954 .table_exists(&table_name[0], &table_name[1], &table_name[2], None)
955 .await
956 .map_err(BoxedError::new)
957 .context(ExternalSnafu)
958 }
959
960 pub(crate) async fn execute_once_serialized(
961 &self,
962 engine: &QueryEngineRef,
963 frontend_client: &Arc<FrontendClient>,
964 max_window_cnt: Option<usize>,
965 ) -> Result<Option<(usize, Duration)>, Error> {
966 let outcome = self
967 .execute_once_serialized_with_outcome(engine, frontend_client, max_window_cnt)
968 .await;
969 outcome.result
970 }
971
972 async fn execute_once_serialized_with_outcome(
975 &self,
976 engine: &QueryEngineRef,
977 frontend_client: &Arc<FrontendClient>,
978 max_window_cnt: Option<usize>,
979 ) -> ExecuteOnceOutcome {
980 let guard = BatchingExecutionGuard::new(self.execution_lock.clone().lock_owned().await);
981 self.execute_once_with_guard(guard, engine, frontend_client, max_window_cnt)
982 .await
983 }
984
985 async fn execute_once_with_guard(
986 &self,
987 guard: BatchingExecutionGuard,
988 engine: &QueryEngineRef,
989 frontend_client: &Arc<FrontendClient>,
990 max_window_cnt: Option<usize>,
991 ) -> ExecuteOnceOutcome {
992 if let Some(execution) = &self.execution {
993 return execution
994 .clone()
995 .execute_once(guard, self, engine, frontend_client, max_window_cnt)
996 .await;
997 }
998 self.execute_once_default_unlocked(engine, frontend_client, max_window_cnt)
999 .await
1000 }
1001
1002 async fn execute_once_default_unlocked(
1004 &self,
1005 engine: &QueryEngineRef,
1006 frontend_client: &Arc<FrontendClient>,
1007 max_window_cnt: Option<usize>,
1008 ) -> ExecuteOnceOutcome {
1009 let new_query = match self.gen_insert_plan_unlocked(engine, max_window_cnt).await {
1010 Ok(new_query) => new_query,
1011 Err(err) => {
1012 return ExecuteOnceOutcome {
1013 new_query: None,
1014 result: Err(err),
1015 };
1016 }
1017 };
1018
1019 if let Some(new_query) = new_query {
1020 debug!("Generate new query: {}", new_query.plan);
1021 let res = self
1022 .execute_logical_plan_unlocked(
1023 engine,
1024 frontend_client,
1025 &new_query.plan,
1026 &new_query.dirty_restore,
1027 &new_query.coverage,
1028 )
1029 .await;
1030 if res.is_err() {
1031 self.handle_executed_query_failure(Some(&new_query));
1032 }
1033 ExecuteOnceOutcome {
1034 new_query: Some(new_query),
1035 result: res,
1036 }
1037 } else {
1038 debug!("Generate no query");
1039 ExecuteOnceOutcome {
1040 new_query: None,
1041 result: Ok(None),
1042 }
1043 }
1044 }
1045
1046 async fn gen_insert_plan_unlocked(
1048 &self,
1049 engine: &QueryEngineRef,
1050 max_window_cnt: Option<usize>,
1051 ) -> Result<Option<PlanInfo>, Error> {
1052 self.gen_insert_plan_with_values_unlocked(engine, max_window_cnt, &BTreeMap::new(), false)
1053 .await
1054 }
1055
1056 pub async fn gen_insert_plan_with_values_unlocked(
1059 &self,
1060 engine: &QueryEngineRef,
1061 max_window_cnt: Option<usize>,
1062 values: &BTreeMap<String, ScalarValue>,
1063 force_full_snapshot: bool,
1064 ) -> Result<Option<PlanInfo>, Error> {
1065 let (table, df_schema) = get_table_info_df_schema(
1066 self.config.catalog_manager.clone(),
1067 self.config.sink_table_name.clone(),
1068 )
1069 .await?;
1070
1071 let table_meta = &table.table_info().meta;
1072 let merge_mode_last_non_null =
1073 is_merge_mode_last_non_null(&table_meta.options.extra_options);
1074 let primary_key_indices = table_meta.primary_key_indices.clone();
1075
1076 let new_query = self
1077 .gen_query_with_time_window_with_values(
1078 engine.clone(),
1079 &table.table_info().meta.schema,
1080 &primary_key_indices,
1081 merge_mode_last_non_null,
1082 max_window_cnt,
1083 values,
1084 force_full_snapshot,
1085 )
1086 .await?;
1087
1088 let Some(new_query) = new_query else {
1089 return Ok(None);
1090 };
1091
1092 let table_columns = df_schema
1095 .columns()
1096 .into_iter()
1097 .map(|c| c.name)
1098 .collect::<BTreeSet<_>>();
1099 for column in new_query.plan.schema().columns() {
1100 if !table_columns.contains(column.name()) {
1101 self.restore_dirty_windows_after_failure(&new_query);
1102 return InvalidQuerySnafu {
1103 reason: format!(
1104 "Column {} not found in sink table with columns {:?}",
1105 column, table_columns
1106 ),
1107 }
1108 .fail();
1109 }
1110 }
1111
1112 let table_provider = Arc::new(DfTableProviderAdapter::new(table));
1113 let table_source = Arc::new(DefaultTableSource::new(table_provider));
1114
1115 let plan = LogicalPlan::Dml(DmlStatement::new(
1117 datafusion_common::TableReference::Full {
1118 catalog: self.config.sink_table_name[0].clone().into(),
1119 schema: self.config.sink_table_name[1].clone().into(),
1120 table: self.config.sink_table_name[2].clone().into(),
1121 },
1122 table_source,
1123 WriteOp::Insert(datafusion_expr::dml::InsertOp::Append),
1124 Arc::new(new_query.plan.clone()),
1125 ));
1126 let insert_into_info = PlanInfo {
1127 plan,
1128 dirty_restore: new_query.dirty_restore,
1129 coverage: new_query.coverage,
1130 };
1131 let insert_into =
1132 match insert_into_info
1133 .plan
1134 .clone()
1135 .recompute_schema()
1136 .context(DatafusionSnafu {
1137 context: "Failed to recompute schema",
1138 }) {
1139 Ok(insert_into) => insert_into,
1140 Err(err) => {
1141 self.restore_dirty_windows_after_failure(&insert_into_info);
1142 return Err(err);
1143 }
1144 };
1145
1146 Ok(Some(PlanInfo {
1147 plan: insert_into,
1148 dirty_restore: insert_into_info.dirty_restore,
1149 coverage: insert_into_info.coverage,
1150 }))
1151 }
1152
1153 pub async fn create_table(
1154 &self,
1155 frontend_client: &Arc<FrontendClient>,
1156 expr: CreateTableExpr,
1157 ) -> Result<(), Error> {
1158 let catalog = &self.config.sink_table_name[0];
1159 let schema = &self.config.sink_table_name[1];
1160 frontend_client
1161 .create(expr.clone(), catalog, schema)
1162 .await?;
1163 Ok(())
1164 }
1165
1166 async fn execute_logical_plan_unlocked(
1168 &self,
1169 engine: &QueryEngineRef,
1170 frontend_client: &Arc<FrontendClient>,
1171 plan: &LogicalPlan,
1172 dirty_restore: &DirtyRestore,
1173 coverage: &QueryCoverage,
1174 ) -> Result<Option<(usize, Duration)>, Error> {
1175 let flow_id = self.config.flow_id;
1176 let Some((res, elapsed)) = self
1177 .execute_plan_unlocked(engine, frontend_client, plan, dirty_restore, coverage)
1178 .await?
1179 else {
1180 return Ok(None);
1181 };
1182
1183 if let Err(err) = &res {
1184 let decision = {
1185 let mut state = self.state.write().unwrap();
1186 let reason = Self::query_failure_reason(err, coverage);
1187 Self::apply_query_failure_to_state(&mut state, elapsed, coverage, reason)
1188 };
1189 if let Some(decision) = decision {
1190 Self::record_checkpoint_decision(flow_id, decision);
1191 }
1192 }
1193
1194 let res = res?;
1195 let (affected_rows, _) = res.output.extract_rows_and_cost();
1196 let decision = {
1197 let mut state = self.state.write().unwrap();
1198 Self::apply_query_result_to_state(&mut state, &res, elapsed, coverage)
1199 };
1200 Self::record_checkpoint_decision(flow_id, decision);
1201 Ok(Some((affected_rows, elapsed)))
1202 }
1203
1204 pub async fn execute_plan_unlocked(
1208 &self,
1209 engine: &QueryEngineRef,
1210 frontend_client: &Arc<FrontendClient>,
1211 plan: &LogicalPlan,
1212 dirty_restore: &DirtyRestore,
1213 coverage: &QueryCoverage,
1214 ) -> Result<Option<(Result<OutputWithMetrics, Error>, Duration)>, Error> {
1215 let instant = Instant::now();
1216 let flow_id = self.config.flow_id;
1217
1218 debug!(
1219 "Executing flow {flow_id}(expire_after={:?} secs) with query {}",
1220 self.config.expire_after, &plan
1221 );
1222
1223 let catalog = &self.config.sink_table_name[0];
1224 let schema = &self.config.sink_table_name[1];
1225
1226 let plan = plan
1228 .clone()
1229 .transform_down_with_subqueries(|p| {
1230 if let LogicalPlan::TableScan(mut table_scan) = p {
1231 let resolved = table_scan.table_name.resolve(catalog, schema);
1232 table_scan.table_name = resolved.into();
1233 Ok(Transformed::yes(LogicalPlan::TableScan(table_scan)))
1234 } else {
1235 Ok(Transformed::no(p))
1236 }
1237 })
1238 .with_context(|_| DatafusionSnafu {
1239 context: format!("Failed to fix table ref in logical plan, plan={:?}", plan),
1240 })?
1241 .data;
1242
1243 let incremental_plan = if coverage.is_incremental_delta() {
1246 self.prepare_plan_for_incremental(engine, &plan).await?
1247 } else {
1248 None
1249 };
1250 let incremental_safe = incremental_plan.is_some();
1251 if coverage.is_incremental_delta() && !incremental_safe {
1252 debug!(
1253 "Flow {flow_id} skipped unsafe incremental delta fallback; \
1254 restored dirty signal instead of executing an unfiltered full snapshot"
1255 );
1256 self.restore_dirty_windows(dirty_restore);
1257 return Ok(None);
1258 }
1259 let plan = incremental_plan.unwrap_or_else(|| plan.clone());
1260 let plan = match &self.execution {
1261 Some(execution) => execution.rewrite_plan(self, plan)?,
1262 None => plan,
1263 };
1264
1265 let extensions = self
1266 .build_flow_query_extensions(incremental_safe, coverage.is_incremental_delta())
1267 .await?;
1268 let frontend_extensions = self.frontend_extensions();
1269 let extension_refs = extensions
1270 .iter()
1271 .map(|(key, value)| (*key, value.as_str()))
1272 .chain(
1273 frontend_extensions
1274 .iter()
1275 .map(|(key, value)| (key.as_str(), value.as_str())),
1276 )
1277 .collect::<Vec<_>>();
1278 let query_mode = if extensions
1279 .iter()
1280 .any(|(key, _)| *key == FLOW_INCREMENTAL_MODE)
1281 {
1282 CheckpointMode::Incremental
1283 } else {
1284 CheckpointMode::FullSnapshot
1285 };
1286 Self::record_query_mode(flow_id, query_mode);
1287 debug!(
1288 "Flow {flow_id} executing batching query with checkpoint_mode={}, extension_count={}",
1289 checkpoint_mode_label(query_mode),
1290 extensions.len()
1291 );
1292
1293 let mut peer_desc = None;
1294 let res = {
1295 let _timer = METRIC_FLOW_BATCHING_ENGINE_QUERY_TIME
1296 .with_label_values(&[flow_id.to_string().as_str()])
1297 .start_timer();
1298
1299 let req = if let Some((insert_to, insert_input_plan)) =
1300 breakup_insert_plan(&plan, catalog, schema)
1301 {
1302 if query_mode == CheckpointMode::FullSnapshot
1303 && matches!(self.config.query_type, QueryType::Sql)
1304 && self.config.flow_eval_interval.is_some()
1305 && self.config.time_window_expr.is_none()
1306 {
1307 match df_plan_to_sql(&insert_input_plan) {
1317 Ok(select_sql) => {
1318 let target_columns = format_insert_target_columns(&insert_input_plan);
1319 let sql = format!(
1320 "INSERT INTO {} ({}) {}",
1321 TableReference::full(
1322 insert_to.catalog_name.as_str(),
1323 insert_to.schema_name.as_str(),
1324 insert_to.table_name.as_str(),
1325 )
1326 .to_quoted_string(),
1327 target_columns,
1328 select_sql
1329 );
1330 api::v1::QueryRequest {
1331 query: Some(api::v1::query_request::Query::Sql(sql)),
1332 }
1333 }
1334 Err(err) => {
1335 debug!(
1336 "Failed to unparse full-snapshot SQL flow {} plan; \
1337 falling back to InsertIntoPlan: {:?}",
1338 flow_id, err
1339 );
1340 encode_insert_plan_request(insert_to, &insert_input_plan)?
1341 }
1342 }
1343 } else {
1344 encode_insert_plan_request(insert_to, &insert_input_plan)?
1345 }
1346 } else {
1347 let message = DFLogicalSubstraitConvertor {}
1348 .encode(&plan, DefaultSerializer)
1349 .context(SubstraitEncodeLogicalPlanSnafu)?;
1350
1351 api::v1::QueryRequest {
1352 query: Some(api::v1::query_request::Query::LogicalPlan(message.to_vec())),
1353 }
1354 };
1355
1356 let snapshot_seqs = coverage.snapshot_seqs();
1357 {
1358 let mut state = self.state.write().unwrap();
1359 state.record_start_time_if_first();
1360 }
1361 frontend_client
1362 .query_with_terminal_metrics(
1363 catalog,
1364 schema,
1365 req,
1366 &extension_refs,
1367 &snapshot_seqs,
1368 &mut peer_desc,
1369 )
1370 .await
1371 };
1372
1373 let elapsed = instant.elapsed();
1374 let peer_label = peer_desc
1375 .as_ref()
1376 .map(ToString::to_string)
1377 .unwrap_or_else(|| PeerDesc::default().to_string());
1378 if let Err(err) = &res {
1379 warn!(
1383 "Failed to execute Flow {flow_id} on frontend {peer_label}, result: {err:?}, elapsed: {:?} with query: {}",
1384 elapsed, &plan
1385 );
1386 }
1387
1388 if elapsed >= self.config.batch_opts.slow_query_threshold {
1390 warn!(
1391 "Flow {flow_id} on frontend {peer_label} executed for {:?} before complete, query: {}",
1392 elapsed, &plan
1393 );
1394 let flow_id = flow_id.to_string();
1395 METRIC_FLOW_BATCHING_ENGINE_SLOW_QUERY
1396 .with_label_values(&[flow_id.as_str(), peer_label.as_str()])
1397 .observe(elapsed.as_secs_f64());
1398 }
1399
1400 match res {
1401 Ok(res) => {
1402 let (affected_rows, _) = res.output.extract_rows_and_cost();
1403 if matches!(&res.output.data, common_query::OutputData::AffectedRows(_)) {
1404 res.metrics.wait_ready().await;
1405 }
1406 if let Some(error) = res.metrics.completion_error() {
1407 warn!("Flow {flow_id} completed with terminal metrics error: {error}");
1408 }
1409 debug!(
1410 "Flow {flow_id} executed, affected_rows: {affected_rows:?}, elapsed: {:?}, watermark: {:?}",
1411 elapsed,
1412 res.region_watermark_map()
1413 );
1414 METRIC_FLOW_ROWS
1415 .with_label_values(&[format!("{}-out-batching", flow_id).as_str()])
1416 .inc_by(affected_rows as _);
1417 Ok(Some((Ok(res), elapsed)))
1418 }
1419 Err(err) => Ok(Some((Err(err), elapsed))),
1420 }
1421 }
1422
1423 pub fn restore_dirty_windows(&self, dirty_restore: &DirtyRestore) {
1427 match dirty_restore {
1428 DirtyRestore::Scoped(filter) => self.restore_scoped_dirty_windows(filter),
1429 DirtyRestore::Unscoped(dirty_windows) => self
1430 .state
1431 .write()
1432 .unwrap()
1433 .dirty_time_windows
1434 .add_dirty_windows(dirty_windows),
1435 }
1436 }
1437
1438 fn restore_dirty_windows_after_failure(&self, query: &PlanInfo) {
1441 self.restore_dirty_windows(&query.dirty_restore);
1442 }
1443
1444 fn restore_scoped_dirty_windows(&self, filter: &FilterExprInfo) {
1447 self.state.write().unwrap().restore_scoped_windows(filter);
1448 }
1449
1450 fn restore_scoped_dirty_windows_on_err<T>(
1453 &self,
1454 filter: &FilterExprInfo,
1455 result: Result<T, Error>,
1456 ) -> Result<T, Error> {
1457 result.inspect_err(|_| {
1458 self.restore_scoped_dirty_windows(filter);
1459 })
1460 }
1461
1462 fn restore_unscoped_dirty_windows(&self, dirty_windows: &DirtyTimeWindows) {
1465 self.state
1466 .write()
1467 .unwrap()
1468 .dirty_time_windows
1469 .add_dirty_windows(dirty_windows);
1470 }
1471
1472 fn restore_unscoped_dirty_windows_on_err<T>(
1475 &self,
1476 dirty_windows: &DirtyTimeWindows,
1477 result: Result<T, Error>,
1478 ) -> Result<T, Error> {
1479 result.inspect_err(|_| {
1480 self.restore_unscoped_dirty_windows(dirty_windows);
1481 })
1482 }
1483
1484 fn drain_dirty_windows_signal(&self) -> (bool, DirtyTimeWindows) {
1487 let mut state = self.state.write().unwrap();
1488 let dirty_windows_to_restore = state.dirty_time_windows.clone();
1489 let is_dirty = !dirty_windows_to_restore.is_empty();
1490 state.dirty_time_windows.clean();
1491 (is_dirty, dirty_windows_to_restore)
1492 }
1493
1494 #[allow(clippy::too_many_arguments)]
1495 async fn gen_unfiltered_plan_info(
1498 &self,
1499 engine: QueryEngineRef,
1500 query_ctx: QueryContextRef,
1501 sink_table_schema: Arc<Schema>,
1502 primary_key_indices: &[usize],
1503 allow_partial: bool,
1504 values: &BTreeMap<String, ScalarValue>,
1505 dirty_windows_to_restore: DirtyTimeWindows,
1506 retention_filter: Option<(&str, Timestamp, &'static str)>,
1507 coverage: QueryCoverage,
1508 ) -> Result<PlanInfo, Error> {
1509 let mut plan = self.restore_unscoped_dirty_windows_on_err(
1510 &dirty_windows_to_restore,
1511 gen_plan_with_matching_schema_and_values(
1512 &self.config.query,
1513 query_ctx,
1514 engine,
1515 sink_table_schema,
1516 primary_key_indices,
1517 allow_partial,
1518 values,
1519 )
1520 .await,
1521 )?;
1522
1523 if let Some((col_name, lower_bound, context)) = retention_filter {
1524 let lower = self.restore_unscoped_dirty_windows_on_err(
1525 &dirty_windows_to_restore,
1526 to_df_literal(lower_bound),
1527 )?;
1528 let retention_filter = col(col_name).gt_eq(lit(lower));
1529 let mut add_filter = AddFilterRewriter::new(retention_filter);
1530 plan = self.restore_unscoped_dirty_windows_on_err(
1531 &dirty_windows_to_restore,
1532 plan.clone()
1533 .rewrite(&mut add_filter)
1534 .with_context(|_| DatafusionSnafu {
1535 context: format!(
1536 "Failed to apply {context} expire_after filter to plan:\n {}\n",
1537 plan
1538 ),
1539 })
1540 .map(|rewrite| rewrite.data),
1541 )?;
1542 }
1543
1544 Ok(PlanInfo {
1545 plan,
1546 dirty_restore: DirtyRestore::Unscoped(dirty_windows_to_restore),
1547 coverage,
1548 })
1549 }
1550
1551 #[allow(clippy::too_many_arguments)]
1552 async fn gen_unfiltered_plan_info_if_dirty(
1555 &self,
1556 engine: QueryEngineRef,
1557 query_ctx: QueryContextRef,
1558 sink_table_schema: Arc<Schema>,
1559 primary_key_indices: &[usize],
1560 allow_partial: bool,
1561 values: &BTreeMap<String, ScalarValue>,
1562 retention_filter: Option<(&str, Timestamp, &'static str)>,
1563 coverage: QueryCoverage,
1564 ) -> Result<Option<PlanInfo>, Error> {
1565 let (is_dirty, dirty_windows_to_restore) = self.drain_dirty_windows_signal();
1566 if !is_dirty {
1567 debug!("Flow id={:?}, no new data, not update", self.config.flow_id);
1568 return Ok(None);
1569 }
1570
1571 self.gen_unfiltered_plan_info(
1572 engine,
1573 query_ctx,
1574 sink_table_schema,
1575 primary_key_indices,
1576 allow_partial,
1577 values,
1578 dirty_windows_to_restore,
1579 retention_filter,
1580 coverage,
1581 )
1582 .await
1583 .map(Some)
1584 }
1585
1586 fn handle_executed_query_failure(&self, query: Option<&PlanInfo>) {
1587 if let Some(query) = query {
1588 self.restore_dirty_windows_after_failure(query);
1589 }
1590 }
1591
1592 pub async fn start_executing_loop(
1600 &self,
1601 engine: QueryEngineRef,
1602 frontend_client: Arc<FrontendClient>,
1603 ) {
1604 if self.config.flow_eval_interval.is_some() {
1605 self.start_scheduled_loop(engine, frontend_client).await;
1606 } else {
1607 self.start_adaptive_loop(engine, frontend_client).await;
1608 }
1609 }
1610
1611 async fn start_scheduled_loop(
1621 &self,
1622 engine: QueryEngineRef,
1623 frontend_client: Arc<FrontendClient>,
1624 ) {
1625 let flow_id_str = self.config.flow_id.to_string();
1626
1627 let schedule = match &self.config.eval_schedule {
1628 Some(s) => s.clone(),
1629 None => {
1630 let eval_interval_secs = self
1631 .config
1632 .flow_eval_interval
1633 .map(|d| d.as_secs() as i64)
1634 .expect("checked by caller");
1635
1636 match EvalSchedule::from_config(Some(eval_interval_secs), None) {
1639 Ok(Some(s)) => s,
1640 Ok(None) => {
1641 warn!(
1642 "Flow {}: EVAL INTERVAL set but no schedule parsed; exiting loop",
1643 flow_id_str
1644 );
1645 return;
1646 }
1647 Err(e) => {
1648 warn!(
1649 "Flow {}: Failed to parse eval schedule: {}; exiting loop",
1650 flow_id_str, e
1651 );
1652 return;
1653 }
1654 }
1655 }
1656 };
1657
1658 let mut cursor_secs =
1663 match initial_schedule_cursor(schedule.start_secs, schedule.interval_secs) {
1664 Ok(cursor) => cursor,
1665 Err(e) => {
1666 warn!(
1667 "Flow {}: invalid eval schedule, exiting loop: {}",
1668 flow_id_str, e
1669 );
1670 return;
1671 }
1672 };
1673
1674 info!(
1675 "Flow {}: entering scheduled loop, interval={}s, start={}, anchor={}, policy={:?}, max_runs={}, max_lag={}s",
1676 flow_id_str,
1677 schedule.interval_secs,
1678 schedule.start_secs,
1679 schedule.anchor_secs,
1680 schedule.missed_tick_policy,
1681 schedule.max_runs,
1682 schedule.max_lag_secs,
1683 );
1684
1685 loop {
1686 if self.is_shutdown_signaled() {
1687 break;
1688 }
1689
1690 let wall_now_secs = wall_clock_unix_secs();
1691
1692 let due = match select_due_scheduled_times(&schedule, cursor_secs, wall_now_secs) {
1693 Ok(d) => d,
1694 Err(e) => {
1695 warn!(
1696 "Flow {}: invalid eval schedule, exiting loop: {}",
1697 flow_id_str, e
1698 );
1699 return;
1700 }
1701 };
1702
1703 if due.scheduled_times_secs.is_empty() {
1704 if due.skipped > 0 {
1705 warn!(
1706 "Flow {}: all {} due scheduled times skipped by max-lag, advancing cursor to wall-clock ({wall_now_secs}) to avoid re-skipping",
1707 flow_id_str, due.skipped
1708 );
1709 cursor_secs = wall_now_secs;
1710 continue;
1711 }
1712
1713 let next = match schedule.next_scheduled_time_after(cursor_secs) {
1715 Ok(next) => next,
1716 Err(e) => {
1717 warn!(
1718 "Flow {}: cannot advance eval schedule past cursor {cursor_secs}: {e}; exiting loop",
1719 flow_id_str
1720 );
1721 return;
1722 }
1723 };
1724 if next <= wall_now_secs {
1725 cursor_secs = wall_now_secs;
1728 continue;
1729 }
1730 let wait_secs = match sleep_delta_secs(next, wall_now_secs) {
1731 Ok(wait_secs) => wait_secs,
1732 Err(e) => {
1733 warn!(
1734 "Flow {}: cannot sleep until next scheduled time {}: {e}; exiting loop",
1735 flow_id_str, next
1736 );
1737 return;
1738 }
1739 };
1740 let wait_dur = Duration::from_secs(wait_secs);
1741 debug!(
1742 "Flow {}: no due scheduled times, sleeping for {}s until next scheduled time at {}",
1743 flow_id_str, wait_secs, next
1744 );
1745 tokio::time::sleep(wait_dur).await;
1746 continue;
1747 }
1748
1749 if due.skipped > 0 {
1750 info!(
1751 "Flow {}: {} due scheduled times, {} skipped (catch-up)",
1752 flow_id_str,
1753 due.scheduled_times_secs.len(),
1754 due.skipped
1755 );
1756 }
1757
1758 for scheduled_time_secs in &due.scheduled_times_secs {
1760 if self.is_shutdown_signaled() {
1761 break;
1762 }
1763
1764 METRIC_FLOW_BATCHING_ENGINE_START_QUERY_CNT
1765 .with_label_values(&[&flow_id_str])
1766 .inc();
1767
1768 let outcome = self
1769 .execute_once_serialized_at_scheduled_time(
1770 &engine,
1771 &frontend_client,
1772 *scheduled_time_secs,
1773 )
1774 .await;
1775
1776 cursor_secs = *scheduled_time_secs;
1778
1779 match outcome.result {
1780 Ok(Some((rows, elapsed))) => {
1781 debug!(
1782 "Flow {}: scheduled time {} completed, rows={}, elapsed={:?}",
1783 flow_id_str, scheduled_time_secs, rows, elapsed
1784 );
1785 }
1786 Ok(None) => {
1787 debug!(
1788 "Flow {}: scheduled time {} produced no query (no dirty signal or no-op)",
1789 flow_id_str, scheduled_time_secs
1790 );
1791 }
1792 Err(err) => {
1793 warn!(
1794 "Flow {}: scheduled time {} failed: {:?}",
1795 flow_id_str, scheduled_time_secs, err
1796 );
1797 METRIC_FLOW_BATCHING_ENGINE_ERROR_CNT
1798 .with_label_values(&[&flow_id_str])
1799 .inc();
1800 }
1804 }
1805 }
1806 }
1807 }
1808
1809 async fn start_adaptive_loop(
1811 &self,
1812 engine: QueryEngineRef,
1813 frontend_client: Arc<FrontendClient>,
1814 ) {
1815 let flow_id_str = self.config.flow_id.to_string();
1816 let mut max_window_cnt = None;
1817 loop {
1818 if self.is_shutdown_signaled() {
1819 break;
1820 }
1821 METRIC_FLOW_BATCHING_ENGINE_START_QUERY_CNT
1822 .with_label_values(&[&flow_id_str])
1823 .inc();
1824
1825 let min_refresh = self.config.batch_opts.experimental_min_refresh_duration;
1826
1827 let outcome = self
1828 .execute_once_serialized_with_outcome(&engine, &frontend_client, max_window_cnt)
1829 .await;
1830
1831 match outcome.result {
1832 Ok(Some(_)) => {
1833 max_window_cnt = max_window_cnt.map(|cnt| {
1834 (cnt + 1).min(self.config.batch_opts.experimental_max_filter_num_per_query)
1835 });
1836
1837 let sleep_until = {
1838 let state = self.state.write().unwrap();
1839
1840 let time_window_size = self
1841 .config
1842 .time_window_expr
1843 .as_ref()
1844 .and_then(|t| *t.time_window_size());
1845
1846 let prefer_short_incremental_cadence = state.checkpoint_mode()
1847 == CheckpointMode::Incremental
1848 && !state.is_incremental_disabled();
1849
1850 state.get_next_start_query_time(
1851 self.config.flow_id,
1852 &time_window_size,
1853 min_refresh,
1854 Some(self.config.batch_opts.query_timeout),
1855 self.config.batch_opts.experimental_max_filter_num_per_query,
1856 prefer_short_incremental_cadence,
1857 )
1858 };
1859
1860 tokio::time::sleep_until(sleep_until).await;
1861 }
1862 Ok(None) => {
1863 debug!(
1864 "Flow id = {:?} found no new data, sleep for {:?} then continue",
1865 self.config.flow_id, min_refresh
1866 );
1867 tokio::time::sleep(min_refresh).await;
1868 continue;
1869 }
1870 Err(err) => {
1871 METRIC_FLOW_BATCHING_ENGINE_ERROR_CNT
1872 .with_label_values(&[&flow_id_str])
1873 .inc();
1874 match outcome.new_query {
1875 Some(query) => {
1876 common_telemetry::error!(err; "Failed to execute query for flow={} with query: {}", self.config.flow_id, query.plan);
1877 max_window_cnt = Some(1);
1878 }
1879 None => {
1880 common_telemetry::error!(err; "Failed to generate query for flow={}", self.config.flow_id)
1881 }
1882 }
1883 tokio::time::sleep(min_refresh).await;
1884 }
1885 }
1886 }
1887 }
1888
1889 fn is_shutdown_signaled(&self) -> bool {
1891 let mut state = self.state.write().unwrap();
1892 match state.shutdown_rx.try_recv() {
1893 Ok(()) | Err(TryRecvError::Closed) => true,
1894 Err(TryRecvError::Empty) => false,
1895 }
1896 }
1897
1898 async fn execute_once_serialized_at_scheduled_time(
1905 &self,
1906 engine: &QueryEngineRef,
1907 frontend_client: &Arc<FrontendClient>,
1908 scheduled_time_secs: i64,
1909 ) -> ExecuteOnceOutcome {
1910 let mut guard = BatchingExecutionGuard::new(self.execution_lock.clone().lock_owned().await);
1911
1912 let scheduled_time_millis = match scheduled_time_millis(scheduled_time_secs) {
1916 Ok(millis) => millis,
1917 Err(e) => {
1918 return ExecuteOnceOutcome {
1919 new_query: None,
1920 result: Err(e),
1921 };
1922 }
1923 };
1924
1925 let old_ctx = {
1928 let mut state = self.state.write().unwrap();
1929 let old = state.query_ctx.clone();
1930 let mut new_ctx = (*old).clone();
1931 new_ctx.set_extension(
1932 query::options::FLOW_SCHEDULED_TIME_MILLIS,
1933 scheduled_time_millis.to_string(),
1934 );
1935 state.query_ctx = Arc::new(new_ctx);
1936 old
1937 };
1938 guard.restore_query_context(self.state.clone(), old_ctx);
1939
1940 self.execute_once_with_guard(guard, engine, frontend_client, None)
1943 .await
1944 }
1945
1946 async fn gen_create_table_expr(
1953 &self,
1954 engine: QueryEngineRef,
1955 ) -> Result<CreateTableExpr, Error> {
1956 let query_ctx = self.state.read().unwrap().query_ctx.clone();
1957 let plan =
1958 sql_to_df_plan(query_ctx.clone(), engine.clone(), &self.config.query, true).await?;
1959 create_table_with_expr(&plan, &self.config.sink_table_name, &self.config.query_type)
1960 }
1961
1962 fn should_use_unfiltered_incremental_delta(&self) -> bool {
1965 let state = self.state.read().unwrap();
1966 state.checkpoint_mode() == CheckpointMode::Incremental
1967 && !state.is_incremental_disabled()
1968 && matches!(self.config.query_type, QueryType::Sql)
1969 }
1970
1971 async fn gen_query_with_time_window(
1973 &self,
1974 engine: QueryEngineRef,
1975 sink_table_schema: &Arc<Schema>,
1976 primary_key_indices: &[usize],
1977 allow_partial: bool,
1978 max_window_cnt: Option<usize>,
1979 ) -> Result<Option<PlanInfo>, Error> {
1980 self.gen_query_with_time_window_with_values(
1981 engine,
1982 sink_table_schema,
1983 primary_key_indices,
1984 allow_partial,
1985 max_window_cnt,
1986 &BTreeMap::new(),
1987 false,
1988 )
1989 .await
1990 }
1991
1992 #[allow(clippy::too_many_arguments)]
1994 async fn gen_query_with_time_window_with_values(
1995 &self,
1996 engine: QueryEngineRef,
1997 sink_table_schema: &Arc<Schema>,
1998 primary_key_indices: &[usize],
1999 allow_partial: bool,
2000 max_window_cnt: Option<usize>,
2001 values: &BTreeMap<String, ScalarValue>,
2002 force_full_snapshot: bool,
2003 ) -> Result<Option<PlanInfo>, Error> {
2004 let query_ctx = self.state.read().unwrap().query_ctx.clone();
2005 let start = SystemTime::now();
2006 let since_the_epoch = start
2007 .duration_since(UNIX_EPOCH)
2008 .expect("Time went backwards");
2009 let low_bound = self
2010 .config
2011 .expire_after
2012 .map(|e| since_the_epoch.as_secs() - e as u64)
2013 .unwrap_or(u64::MIN);
2014
2015 let low_bound = Timestamp::new_second(low_bound as i64);
2016
2017 let expire_time_window_bound = self
2018 .config
2019 .time_window_expr
2020 .as_ref()
2021 .map(|expr| expr.eval(low_bound))
2022 .transpose()?;
2023
2024 let (expire_lower_bound, expire_upper_bound) = match (
2025 expire_time_window_bound,
2026 &self.config.query_type,
2027 ) {
2028 (Some((Some(l), Some(u))), QueryType::Sql) => (l, u),
2029 (None, QueryType::Sql) if self.config.flow_eval_interval.is_none() => {
2030 return UnexpectedSnafu {
2031 reason: format!(
2032 "Flow id={} reached execution without a time-window expression or EVAL INTERVAL; create-flow validation should have rejected it",
2033 self.config.flow_id
2034 ),
2035 }
2036 .fail();
2037 }
2038 _ => {
2039 let (_, dirty_windows_to_restore) = self.drain_dirty_windows_signal();
2045
2046 let plan_info = self
2047 .gen_unfiltered_plan_info(
2048 engine,
2049 query_ctx,
2050 sink_table_schema.clone(),
2051 primary_key_indices,
2052 allow_partial,
2053 values,
2054 dirty_windows_to_restore,
2055 None,
2056 QueryCoverage::UnfilteredFull,
2057 )
2058 .await?;
2059
2060 return Ok(Some(plan_info));
2061 }
2062 };
2063
2064 debug!(
2065 "Flow id = {:?}, found time window: precise_lower_bound={:?}, precise_upper_bound={:?} with dirty time windows: {:?}",
2066 self.config.flow_id,
2067 expire_lower_bound,
2068 expire_upper_bound,
2069 self.state.read().unwrap().dirty_time_windows
2070 );
2071 let window_size = expire_upper_bound
2072 .sub(&expire_lower_bound)
2073 .with_context(|| UnexpectedSnafu {
2074 reason: format!(
2075 "Can't get window size from {expire_upper_bound:?} - {expire_lower_bound:?}"
2076 ),
2077 })?;
2078 let col_name = self
2079 .config
2080 .time_window_expr
2081 .as_ref()
2082 .map(|expr| expr.column_name.clone())
2083 .with_context(|| UnexpectedSnafu {
2084 reason: format!(
2085 "Flow id={:?}, Failed to get column name from time window expr",
2086 self.config.flow_id
2087 ),
2088 })?;
2089
2090 if force_full_snapshot {
2091 let retention_filter = self.config.expire_after.map(|_| {
2095 (
2096 col_name.as_str(),
2097 expire_lower_bound,
2098 "forced full snapshot",
2099 )
2100 });
2101 let (_, dirty_windows_to_restore) = self.drain_dirty_windows_signal();
2102 return self
2103 .gen_unfiltered_plan_info(
2104 engine,
2105 query_ctx,
2106 sink_table_schema.clone(),
2107 primary_key_indices,
2108 allow_partial,
2109 values,
2110 dirty_windows_to_restore,
2111 retention_filter,
2112 QueryCoverage::UnfilteredFull,
2113 )
2114 .await
2115 .map(Some);
2116 }
2117
2118 if self.should_use_unfiltered_incremental_delta() {
2119 let retention_filter = self
2129 .config
2130 .expire_after
2131 .map(|_| (col_name.as_str(), expire_lower_bound, "incremental"));
2132 return self
2133 .gen_unfiltered_plan_info_if_dirty(
2134 engine,
2135 query_ctx,
2136 sink_table_schema.clone(),
2137 primary_key_indices,
2138 allow_partial,
2139 values,
2140 retention_filter,
2141 QueryCoverage::IncrementalDelta,
2142 )
2143 .await;
2144 }
2145
2146 let (expr, coverage) = {
2147 let mut state = self.state.write().unwrap();
2148 let window_cnt = max_window_cnt
2149 .unwrap_or(self.config.batch_opts.experimental_max_filter_num_per_query);
2150 let expr = state.gen_scoped_filter_exprs(
2151 &col_name,
2152 Some(expire_lower_bound),
2153 window_size,
2154 window_cnt,
2155 self.config.flow_id,
2156 Some(self),
2157 )?;
2158 let repair_high = state
2159 .pending_fenced_repair()
2160 .map(|repair| repair.high().clone());
2161 let coverage = if let Some(high) = repair_high {
2162 QueryCoverage::FencedRepairChunk { high }
2163 } else {
2164 QueryCoverage::ScopedBaseRepair
2165 };
2166 (expr, coverage)
2167 };
2168
2169 let Some(expr) = expr else {
2170 debug!("Flow id={:?}, no new data, not update", self.config.flow_id);
2172 return Ok(None);
2173 };
2174
2175 let filter_sql = expr_to_sql(&expr.expr)
2176 .map(|sql| sql.to_string())
2177 .unwrap_or_else(|err| format!("<failed to format filter expr: {err}>"));
2178
2179 debug!(
2180 "Flow id={:?}, Generated filter expr: {:?}",
2181 self.config.flow_id, filter_sql
2182 );
2183
2184 let mut add_filter = AddFilterRewriter::new(expr.expr.clone());
2185 let mut add_auto_column = ColumnMatcherRewriter::new_with_values(
2186 sink_table_schema.clone(),
2187 primary_key_indices.to_vec(),
2188 allow_partial,
2189 values.clone(),
2190 );
2191
2192 let plan = self.restore_scoped_dirty_windows_on_err(
2193 &expr,
2194 sql_to_df_plan(query_ctx.clone(), engine.clone(), &self.config.query, false).await,
2195 )?;
2196 let rewrite = self.restore_scoped_dirty_windows_on_err(
2197 &expr,
2198 plan.clone()
2199 .rewrite(&mut add_filter)
2200 .and_then(|p| p.data.rewrite(&mut add_auto_column))
2201 .with_context(|_| DatafusionSnafu {
2202 context: format!("Failed to rewrite plan:\n {}\n", plan),
2203 })
2204 .map(|rewrite| rewrite.data),
2205 )?;
2206 let new_plan = self.restore_scoped_dirty_windows_on_err(
2208 &expr,
2209 apply_df_optimizer(rewrite, &query_ctx).await,
2210 )?;
2211
2212 let info = PlanInfo {
2213 plan: new_plan.clone(),
2214 dirty_restore: DirtyRestore::Scoped(expr),
2215 coverage,
2216 };
2217
2218 Ok(Some(info))
2219 }
2220}
2221
2222#[cfg(test)]
2223mod test;