1use std::collections::BTreeMap;
18use std::sync::Arc;
19#[cfg(test)]
20use std::sync::atomic::AtomicUsize;
21use std::sync::atomic::{AtomicBool, Ordering};
22
23use common_base::memory_limit::MemoryLimit;
24use common_config::Configurable;
25use common_error::ext::BoxedError;
26use common_meta::distributed_time_constants::BASE_HEARTBEAT_INTERVAL;
27use common_meta::key::TableMetadataManagerRef;
28use common_options::memory::MemoryOptions;
29use common_recordbatch::map_dictionary_to_values_data_type;
30use common_stat::get_total_cpu_cores;
31use common_telemetry::logging::{LoggingOptions, TracingOptions};
32use common_telemetry::{error, info};
33use datafusion_common::TableReference;
34use datafusion_common::tree_node::TreeNode;
35use datafusion_expr::logical_plan::Distinct;
36use datafusion_expr::{Expr, LogicalPlan};
37use datatypes::schema::{ColumnSchema, SchemaRef};
38use itertools::Itertools;
39use meta_client::MetaClientOptions;
40use query::QueryEngine;
41use query::options::QueryOptions;
42use serde::{Deserialize, Serialize};
43use servers::grpc::GrpcOptions;
44use servers::http::HttpOptions;
45use snafu::{IntoError, OptionExt, ResultExt, ensure};
46use store_api::storage::{ConcreteDataType, RegionId};
47use tokio::sync::RwLock;
48use tokio::time::Instant;
49
50use crate::adapter::stateless::StatelessFlow;
51use crate::adapter::table_source::ManagedTableSource;
52use crate::adapter::util::{
53 relation_desc_to_column_schemas_with_fallback, table_info_value_to_relation_desc,
54};
55use crate::batching_mode::BatchingModeOptions;
56use crate::batching_mode::frontend_client::FrontendClient;
57use crate::batching_mode::utils::sql_to_df_plan;
58use crate::error::{
59 DatafusionSnafu, Error, ExternalSnafu, FlowNotFoundSnafu, InsertIntoFlowSnafu, InternalSnafu,
60 InvalidQuerySnafu, UnexpectedSnafu,
61};
62use crate::repr::{ColumnType, DiffRow, RelationDesc, Row};
63use crate::{CreateFlowArgs, FlowId, TableName};
64
65pub(crate) mod flownode_impl;
66pub(crate) mod stateless;
67pub(crate) mod table_source;
68#[cfg(test)]
69mod tests;
70pub(crate) mod util;
71
72fn output_column_schemas(
76 plan: &LogicalPlan,
77 source_schema: &SchemaRef,
78) -> Result<(Vec<ColumnSchema>, Vec<Option<usize>>), Error> {
79 let expressions = match plan {
82 LogicalPlan::Projection(projection) => Some(&projection.expr),
83 LogicalPlan::Distinct(Distinct::All(input)) => match input.as_ref() {
84 LogicalPlan::Projection(projection) => Some(&projection.expr),
85 _ => None,
86 },
87 _ => None,
88 };
89 let is_pass_through = matches!(plan, LogicalPlan::TableScan(_) | LogicalPlan::Filter(_))
90 || matches!(plan, LogicalPlan::Distinct(Distinct::All(input)) if matches!(input.as_ref(), LogicalPlan::TableScan(_) | LogicalPlan::Filter(_)));
91 let mut lineage = Vec::new();
92 let columns = plan
93 .schema()
94 .fields()
95 .iter()
96 .enumerate()
97 .map(|(idx, field)| {
98 let mut column = ColumnSchema::try_from(field.as_ref())
99 .map_err(BoxedError::new)
100 .context(ExternalSnafu)?;
101 let source_index = expressions
102 .and_then(|exprs| {
103 let expr = exprs.get(idx)?;
104 let expr = match expr {
105 Expr::Column(column) => Some(column),
106 Expr::Alias(alias) => match alias.expr.as_ref() {
107 Expr::Column(column) => Some(column),
108 _ => None,
109 },
110 _ => None,
111 }?;
112 source_schema.column_index_by_name(&expr.name)
113 })
114 .or_else(|| {
115 is_pass_through
116 .then(|| source_schema.column_index_by_name(field.name()))
117 .flatten()
118 });
119 if let Some(source_index) = source_index {
120 let mut source = source_schema.column_schemas()[source_index].clone();
121 source.name = field.name().clone();
122 source.data_type = map_dictionary_to_values_data_type(&source.data_type);
123 column = source;
124 } else {
125 column.data_type = map_dictionary_to_values_data_type(&column.data_type);
126 if expressions.is_some() {
130 column = column.with_time_index(false);
131 }
132 }
133 lineage.push(source_index);
134 Ok(column)
135 })
136 .collect::<Result<Vec<_>, Error>>()?;
137 Ok((columns, lineage))
138}
139
140fn relation_desc_from_output(
141 columns: &[ColumnSchema],
142 lineage: &[Option<usize>],
143 source_primary_key_indices: &[usize],
144) -> RelationDesc {
145 let keys = source_primary_key_indices
146 .iter()
147 .filter_map(|source_index| {
148 lineage
149 .iter()
150 .position(|index| index == &Some(*source_index))
151 })
152 .collect_vec();
153 let time_index = columns.iter().position(ColumnSchema::is_time_index);
154 RelationDesc {
155 typ: crate::repr::RelationType {
156 column_types: columns
157 .iter()
158 .map(|column| ColumnType::new(column.data_type.clone(), column.is_nullable()))
159 .collect(),
160 keys: if keys.is_empty() {
161 vec![]
162 } else {
163 vec![crate::repr::Key::from(keys)]
164 },
165 time_index,
166 auto_columns: vec![],
167 },
168 names: columns
169 .iter()
170 .map(|column| Some(column.name.clone()))
171 .collect(),
172 }
173}
174
175fn default_num_workers() -> usize {
176 get_total_cpu_cores().div_ceil(2)
177}
178
179pub(crate) fn is_explicit_source_timestamp_compatibility(
185 output_schema: &[ColumnSchema],
186 output_lineage: &[Option<usize>],
187 sink_schema: &[ColumnSchema],
188 source_schema: &SchemaRef,
189) -> bool {
190 if sink_schema.len() != output_schema.len() + 1 || output_lineage.len() != output_schema.len() {
191 return false;
192 }
193 if !output_schema
194 .iter()
195 .zip(&sink_schema[..output_schema.len()])
196 .all(|(output, sink)| output.data_type == sink.data_type)
197 {
198 return false;
199 }
200
201 let Some(source_timestamp_index) = source_schema.timestamp_index() else {
202 return false;
203 };
204 let sink_timestamp = &sink_schema[output_schema.len()];
205 sink_timestamp.data_type == source_schema.column_schemas()[source_timestamp_index].data_type
206 && sink_timestamp.data_type.is_timestamp()
207 && sink_timestamp.is_time_index()
208 && sink_timestamp.default_constraint().is_some()
209 && sink_timestamp.name != AUTO_CREATED_UPDATE_AT_TS_COL
210 && sink_timestamp.name != AUTO_CREATED_PLACEHOLDER_TS_COL
211 && !output_lineage.contains(&Some(source_timestamp_index))
212}
213
214pub const AUTO_CREATED_PLACEHOLDER_TS_COL: &str = "__ts_placeholder";
215pub const AUTO_CREATED_UPDATE_AT_TS_COL: &str = "update_at";
216
217pub(crate) fn resolve_sink_layout(
223 output_schema: &[ColumnSchema],
224 sink_schema: &[ColumnSchema],
225) -> Result<Vec<ColumnSchema>, Error> {
226 ensure!(
227 sink_schema.len() >= output_schema.len() && sink_schema.len() - output_schema.len() <= 2,
228 InvalidQuerySnafu {
229 reason: format!(
230 "Flow output has {} columns, but sink has {} columns; only zero, one, or two trailing auto columns are supported",
231 output_schema.len(),
232 sink_schema.len()
233 )
234 }
235 );
236 let suffix_len = sink_schema.len() - output_schema.len();
237 for (idx, (output, sink)) in output_schema.iter().zip(sink_schema.iter()).enumerate() {
238 ensure!(
239 output.data_type == sink.data_type,
240 InvalidQuerySnafu {
241 reason: format!(
242 "Flow output column {idx} has type {:?}, but sink column {} has type {:?}",
243 output.data_type, sink.name, sink.data_type
244 )
245 }
246 );
247 }
248 let suffix = &sink_schema[output_schema.len()..];
249 match suffix_len {
250 0 => Ok(vec![]),
251 1 => {
252 let column = &suffix[0];
253 ensure!(
254 column.name == AUTO_CREATED_UPDATE_AT_TS_COL && column.data_type.is_timestamp(),
255 InvalidQuerySnafu {
256 reason: format!(
257 "The trailing sink column must be timestamp {}",
258 AUTO_CREATED_UPDATE_AT_TS_COL
259 )
260 }
261 );
262 Ok(suffix.to_vec())
263 }
264 2 => {
265 let update_at = &suffix[0];
266 let placeholder = &suffix[1];
267 ensure!(
268 update_at.name == AUTO_CREATED_UPDATE_AT_TS_COL
269 && update_at.data_type.is_timestamp()
270 && placeholder.name == AUTO_CREATED_PLACEHOLDER_TS_COL
271 && placeholder.data_type.is_timestamp()
272 && placeholder.is_time_index(),
273 InvalidQuerySnafu {
274 reason: "The two trailing sink columns must be timestamp update_at followed by timestamp time-index __ts_placeholder".to_string()
275 }
276 );
277 Ok(suffix.to_vec())
278 }
279 _ => unreachable!(),
280 }
281}
282
283pub(crate) fn sink_output_column_count(sink_schema: &[ColumnSchema]) -> Result<usize, Error> {
286 let mut count = sink_schema.len();
287 if sink_schema
288 .last()
289 .is_some_and(|column| column.name == AUTO_CREATED_PLACEHOLDER_TS_COL)
290 {
291 let placeholder = &sink_schema[count - 1];
292 ensure!(
293 placeholder.data_type.is_timestamp() && placeholder.is_time_index(),
294 InvalidQuerySnafu {
295 reason: format!(
296 "Auto-created sink column {} must be a timestamp time index",
297 AUTO_CREATED_PLACEHOLDER_TS_COL
298 )
299 }
300 );
301 count -= 1;
302 }
303 if sink_schema
304 .get(count.saturating_sub(1))
305 .is_some_and(|column| column.name == AUTO_CREATED_UPDATE_AT_TS_COL)
306 {
307 ensure!(
308 sink_schema[count - 1].data_type.is_timestamp(),
309 InvalidQuerySnafu {
310 reason: format!(
311 "Auto-created sink column {} must be a timestamp",
312 AUTO_CREATED_UPDATE_AT_TS_COL
313 )
314 }
315 );
316 count -= 1;
317 }
318 Ok(count)
319}
320
321pub(crate) fn validate_sink_layout(
322 output_schema: &[ColumnSchema],
323 sink_schema: &[ColumnSchema],
324) -> Result<(), Error> {
325 resolve_sink_layout(output_schema, sink_schema).map(|_| ())
326}
327
328pub(crate) fn validate_auto_column_names(output_schema: &[ColumnSchema]) -> Result<(), Error> {
329 for column in output_schema {
330 ensure!(
331 column.name != AUTO_CREATED_UPDATE_AT_TS_COL
332 && column.name != AUTO_CREATED_PLACEHOLDER_TS_COL,
333 InvalidQuerySnafu {
334 reason: format!(
335 "Flow output column {} is reserved for an auto-created sink column",
336 column.name
337 )
338 }
339 );
340 }
341 Ok(())
342}
343
344pub(crate) fn validate_sink_layout_with_suffix(
345 output_schema: &[ColumnSchema],
346 sink_schema: &[ColumnSchema],
347 suffix: &[ColumnSchema],
348) -> Result<(), Error> {
349 ensure!(
350 sink_schema.len() == output_schema.len() + suffix.len()
351 && sink_schema[output_schema.len()..] == *suffix,
352 InvalidQuerySnafu {
353 reason: "Stored sink auto-column layout no longer matches the sink schema".to_string()
354 }
355 );
356 for (idx, (output, sink)) in output_schema
357 .iter()
358 .zip(&sink_schema[..output_schema.len()])
359 .enumerate()
360 {
361 ensure!(
362 output.data_type == sink.data_type,
363 InvalidQuerySnafu {
364 reason: format!(
365 "Flow output column {idx} has type {:?}, but sink column {} has type {:?}",
366 output.data_type, sink.name, sink.data_type
367 )
368 }
369 );
370 }
371 Ok(())
372}
373
374#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
375#[serde(default)]
376pub struct FlowConfig {
377 #[deprecated(note = "flow workers have been removed; this field is ignored")]
379 #[serde(default = "default_num_workers")]
380 pub num_workers: usize,
381 pub batching_mode: BatchingModeOptions,
382}
383
384#[allow(deprecated)]
385impl Default for FlowConfig {
386 fn default() -> Self {
387 Self {
388 num_workers: get_total_cpu_cores().div_ceil(2),
389 batching_mode: BatchingModeOptions::default(),
390 }
391 }
392}
393
394#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
395#[serde(default)]
396pub struct FlownodeOptions {
397 pub node_id: Option<u64>,
398 pub flow: FlowConfig,
399 pub grpc: GrpcOptions,
400 pub http: HttpOptions,
401 pub meta_client: Option<MetaClientOptions>,
402 pub logging: LoggingOptions,
403 pub tracing: TracingOptions,
404 pub query: QueryOptions,
405 pub memory: MemoryOptions,
406}
407
408impl Default for FlownodeOptions {
409 fn default() -> Self {
410 Self {
411 node_id: None,
412 flow: FlowConfig::default(),
413 grpc: GrpcOptions::default().with_bind_addr("127.0.0.1:3004"),
414 http: HttpOptions::default(),
415 meta_client: None,
416 logging: LoggingOptions::default(),
417 tracing: TracingOptions::default(),
418 query: QueryOptions {
419 parallelism: 1,
420 allow_query_fallback: false,
421 memory_pool_size: MemoryLimit::default(),
422 enable_per_region_metrics: false,
423 ..Default::default()
424 },
425 memory: MemoryOptions::default(),
426 }
427 }
428}
429
430impl Configurable for FlownodeOptions {
431 fn validate_sanitize(&mut self) -> common_config::error::Result<()> {
432 Ok(())
433 }
434}
435
436pub type FlowStreamingEngineRef = Arc<StreamingEngine>;
437
438#[derive(Default)]
439struct StatelessFlowRuntime {
440 flow: Option<Arc<StatelessFlow>>,
441 failed_rebuild: Option<(u32, Instant)>,
442}
443
444struct StatelessFlowSlot {
445 runtime: Arc<RwLock<StatelessFlowRuntime>>,
446 active: std::sync::atomic::AtomicBool,
448 #[cfg(test)]
449 rebuild_attempts: AtomicUsize,
450}
451
452fn validate_captured_slot(
453 slot: &StatelessFlowSlot,
454 current_source_table_id: Option<table::metadata::TableId>,
455 expected_table_id: table::metadata::TableId,
456 flow_id: FlowId,
457) -> Result<(), Error> {
458 ensure!(
459 slot.active.load(Ordering::Acquire),
460 FlowNotFoundSnafu { id: flow_id }
461 );
462 let current_source_table_id =
463 current_source_table_id.context(FlowNotFoundSnafu { id: flow_id })?;
464 ensure!(
465 current_source_table_id == expected_table_id,
466 InvalidQuerySnafu {
467 reason: format!("Flow {flow_id} source table changed while it was selected")
468 }
469 );
470 Ok(())
471}
472
473pub struct StreamingEngine {
474 pub query_engine: Arc<dyn QueryEngine>,
475 pub frontend_client: Arc<FrontendClient>,
476 table_info_source: ManagedTableSource,
477 stateless_flows: RwLock<BTreeMap<FlowId, Arc<StatelessFlowSlot>>>,
478 pub node_id: Option<u32>,
479}
480
481impl StreamingEngine {
482 pub fn new(
483 node_id: Option<u32>,
484 query_engine: Arc<dyn QueryEngine>,
485 table_meta: TableMetadataManagerRef,
486 frontend_client: Arc<FrontendClient>,
487 ) -> Self {
488 let table_info_source = ManagedTableSource::new(
489 table_meta.table_info_manager().clone(),
490 table_meta.table_name_manager().clone(),
491 );
492 Self {
493 query_engine,
494 frontend_client,
495 table_info_source,
496 stateless_flows: Default::default(),
497 node_id,
498 }
499 }
500
501 pub async fn handle_write_request(
502 &self,
503 region_id: RegionId,
504 rows: Vec<DiffRow>,
505 batch_datatypes: &[ConcreteDataType],
506 source_schema_version: u32,
507 ) -> Result<(), Error> {
508 let table_id = region_id.table_id();
509 let flow_ids = self.flow_ids_for_table(table_id).await;
510 let mut failed_flow_ids = Vec::new();
511 let mut first_error = None;
512 for (flow_id, slot) in flow_ids {
513 let result = self
514 .execute_flow(
515 flow_id,
516 slot,
517 table_id,
518 rows.clone(),
519 batch_datatypes,
520 source_schema_version,
521 )
522 .await;
523 if let Err(err) = result {
524 error!(err; "Failed to insert into flow={}, region_id={}", flow_id, region_id);
525 failed_flow_ids.push(flow_id);
526 if first_error.is_none() {
527 first_error = Some(BoxedError::new(err));
528 }
529 }
530 }
531 match first_error {
532 Some(source) => Err(InsertIntoFlowSnafu {
533 region_id: u64::from(region_id),
534 flow_ids: failed_flow_ids,
535 }
536 .into_error(source)),
537 None => Ok(()),
538 }
539 }
540
541 async fn flow_ids_for_table(
542 &self,
543 table_id: table::metadata::TableId,
544 ) -> Vec<(FlowId, Arc<StatelessFlowSlot>)> {
545 let slots = self
546 .stateless_flows
547 .read()
548 .await
549 .iter()
550 .map(|(id, slot)| (*id, Arc::clone(slot)))
551 .collect::<Vec<_>>();
552 let mut flows = Vec::new();
553 for (id, slot) in slots {
554 if slot
555 .runtime
556 .read()
557 .await
558 .flow
559 .as_ref()
560 .is_some_and(|flow| flow.source_table_id == table_id)
561 {
562 flows.push((id, slot));
563 }
564 }
565 flows
566 }
567
568 async fn execute_flow(
569 &self,
570 flow_id: FlowId,
571 slot: Arc<StatelessFlowSlot>,
572 expected_table_id: table::metadata::TableId,
573 rows: Vec<DiffRow>,
574 batch_datatypes: &[ConcreteDataType],
575 source_schema_version: u32,
576 ) -> Result<usize, Error> {
577 let mut guard = slot.runtime.clone().read_owned().await;
579 let current = guard
580 .flow
581 .as_ref()
582 .context(FlowNotFoundSnafu { id: flow_id })?;
583 validate_captured_slot(
584 &slot,
585 Some(current.source_table_id),
586 expected_table_id,
587 flow_id,
588 )?;
589 if current.source_schema_version != source_schema_version {
590 if guard
591 .failed_rebuild
592 .as_ref()
593 .is_some_and(|(version, until)| {
594 *version == source_schema_version && *until > Instant::now()
595 })
596 {
597 return InvalidQuerySnafu {
598 reason: format!(
599 "Flow {flow_id} schema rebuild for source version {source_schema_version} is cooling down"
600 ),
601 }
602 .fail();
603 }
604 drop(guard);
605 let mut published = slot.runtime.clone().write_owned().await;
608 let current = published
609 .flow
610 .as_ref()
611 .context(FlowNotFoundSnafu { id: flow_id })?;
612 validate_captured_slot(
613 &slot,
614 Some(current.source_table_id),
615 expected_table_id,
616 flow_id,
617 )?;
618 if current.source_schema_version != source_schema_version {
619 if published
620 .failed_rebuild
621 .as_ref()
622 .is_some_and(|(version, until)| {
623 *version == source_schema_version && *until > Instant::now()
624 })
625 {
626 return InvalidQuerySnafu {
627 reason: format!(
628 "Flow {flow_id} schema rebuild for source version {source_schema_version} is cooling down"
629 ),
630 }
631 .fail();
632 }
633 #[cfg(test)]
634 slot.rebuild_attempts.fetch_add(1, Ordering::Relaxed);
635 let replacement = self
636 .build_stateless_flow(¤t.create_args, false)
637 .await
638 .and_then(|replacement| {
639 ensure!(
640 replacement.source_table_id == expected_table_id
641 && replacement.source_schema_version == source_schema_version,
642 InvalidQuerySnafu {
643 reason: format!(
644 "Source schema changed while rebuilding flow {flow_id}"
645 )
646 }
647 );
648 Ok(replacement)
649 });
650 match replacement {
651 Ok(replacement) => {
652 published.flow = Some(Arc::new(replacement));
653 published.failed_rebuild = None;
654 }
655 Err(error) => {
656 published.failed_rebuild = Some((
659 source_schema_version,
660 Instant::now() + BASE_HEARTBEAT_INTERVAL,
661 ));
662 return Err(error);
663 }
664 }
665 }
666 guard = tokio::sync::OwnedRwLockWriteGuard::downgrade(published);
668 }
669 let flow = guard
670 .flow
671 .as_ref()
672 .context(FlowNotFoundSnafu { id: flow_id })?;
673 let latest = self
676 .table_info_source
677 .get_table_info_value(&flow.source_table_id)
678 .await?
679 .context(UnexpectedSnafu {
680 reason: "Source table metadata is missing",
681 })?
682 .table_info
683 .meta
684 .schema
685 .version();
686 ensure!(
687 source_schema_version == latest && flow.source_schema_version == latest,
688 InvalidQuerySnafu {
689 reason: format!("Source schema version changed before flow {flow_id} execution")
690 }
691 );
692 stateless::execute(
693 flow,
694 &rows,
695 batch_datatypes,
696 &self.query_engine,
697 &self.frontend_client,
698 latest,
699 )
700 .await
701 }
702
703 pub async fn remove_flow_inner(&self, flow_id: FlowId) -> Result<(), Error> {
704 let slot = self
707 .stateless_flows
708 .write()
709 .await
710 .remove(&flow_id)
711 .context(FlowNotFoundSnafu { id: flow_id })?;
712 slot.active.store(false, Ordering::Release);
713 let mut runtime = slot.runtime.write().await;
714 runtime.flow.take();
715 runtime.failed_rebuild = None;
716 Ok(())
717 }
718
719 async fn publish_initial_flow(
720 &self,
721 slot: Arc<StatelessFlowSlot>,
722 flow: Arc<StatelessFlow>,
723 create_if_not_exists: bool,
724 or_replace: bool,
725 ) -> Result<bool, Error> {
726 let mut runtime = slot.runtime.write().await;
727 ensure!(
728 slot.active.load(Ordering::Acquire),
729 FlowNotFoundSnafu {
730 id: flow.create_args.flow_id
731 }
732 );
733 if runtime.flow.is_some() && !or_replace {
734 if create_if_not_exists {
735 return Ok(false);
736 }
737 return crate::error::FlowAlreadyExistSnafu {
738 id: flow.create_args.flow_id,
739 }
740 .fail();
741 }
742 let latest = self
743 .table_info_source
744 .get_table_info_value(&flow.source_table_id)
745 .await?
746 .context(UnexpectedSnafu {
747 reason: "Source table metadata is missing",
748 })?
749 .table_info
750 .meta
751 .schema
752 .version();
753 ensure!(
754 latest == flow.source_schema_version,
755 InvalidQuerySnafu {
756 reason: format!(
757 "Source schema changed while building flow: built version {}, current version {latest}",
758 flow.source_schema_version
759 )
760 }
761 );
762 runtime.flow = Some(flow);
763 runtime.failed_rebuild = None;
764 Ok(true)
765 }
766
767 pub async fn create_flow_inner(&self, args: CreateFlowArgs) -> Result<Option<FlowId>, Error> {
768 let flow_id = args.flow_id;
769 let flow = Arc::new(self.build_stateless_flow(&args, true).await?);
770 let slot = {
771 let mut slots = self.stateless_flows.write().await;
772 if let Some(slot) = slots.get(&flow_id) {
773 Arc::clone(slot)
774 } else {
775 let slot = Arc::new(StatelessFlowSlot {
776 runtime: Arc::new(RwLock::new(StatelessFlowRuntime::default())),
777 active: AtomicBool::new(true),
778 #[cfg(test)]
779 rebuild_attempts: AtomicUsize::new(0),
780 });
781 slots.insert(flow_id, Arc::clone(&slot));
782 slot
783 }
784 };
785 let published = self
786 .publish_initial_flow(
787 slot.clone(),
788 flow,
789 args.create_if_not_exists,
790 args.or_replace,
791 )
792 .await?;
793 ensure!(
794 self.stateless_flows
795 .read()
796 .await
797 .get(&flow_id)
798 .is_some_and(|current| Arc::ptr_eq(current, &slot)),
799 FlowNotFoundSnafu { id: flow_id }
800 );
801 if published {
802 info!("Successfully create flow with id={flow_id}");
803 Ok(Some(flow_id))
804 } else {
805 Ok(None)
806 }
807 }
808
809 async fn build_stateless_flow(
810 &self,
811 args: &CreateFlowArgs,
812 create_sink: bool,
813 ) -> Result<StatelessFlow, Error> {
814 let CreateFlowArgs {
815 flow_id,
816 sink_table_name,
817 source_table_ids,
818 sql,
819 query_ctx,
820 ..
821 } = args;
822 ensure!(
823 source_table_ids.len() == 1,
824 InvalidQuerySnafu {
825 reason: "Stateless streaming flow does not support multiple source tables",
826 }
827 );
828
829 let query_ctx = query_ctx.clone().map(Arc::new).context(UnexpectedSnafu {
830 reason: "Query context is missing",
831 })?;
832 let source_table_id = source_table_ids[0];
833 let source_table_name = self
834 .table_info_source
835 .get_table_name(&source_table_id)
836 .await?;
837 let source_table_info = self
838 .table_info_source
839 .get_table_info_value(&source_table_id)
840 .await?
841 .context(UnexpectedSnafu {
842 reason: "Source table metadata is missing",
843 })?;
844 let flow_plan =
845 sql_to_df_plan(query_ctx.clone(), self.query_engine.clone(), sql, true).await?;
846 stateless::validate_plan(&flow_plan)?;
847 let source_meta = source_table_info.table_info.meta;
848 let source_schema = source_meta.schema;
849 let source_primary_key_indices = source_meta.primary_key_indices;
850 stateless::validate_source_scan(&flow_plan, source_table_id, &source_schema)?;
851 let (inferred_schema, lineage) = output_column_schemas(&flow_plan, &source_schema)?;
852 let inferred_relation =
853 relation_desc_from_output(&inferred_schema, &lineage, &source_primary_key_indices);
854 let sink_exists = self.fetch_table_pk_schema(sink_table_name).await?.is_some();
855 if !sink_exists {
856 validate_auto_column_names(&inferred_schema)?;
857 }
858 if !sink_exists
859 && create_sink
860 && !self
861 .create_table_from_relation(
862 &format!("flow-id={flow_id}"),
863 sink_table_name,
864 &inferred_relation,
865 )
866 .await?
867 {
868 return UnexpectedSnafu {
869 reason: format!("Failed to create table {sink_table_name:?}"),
870 }
871 .fail();
872 }
873 ensure!(
874 sink_exists || create_sink,
875 UnexpectedSnafu {
876 reason: format!("Sink table metadata is missing: {sink_table_name:?}"),
877 }
878 );
879
880 let (sink_primary_keys, _, sink_schema) = self
883 .fetch_table_pk_schema(sink_table_name)
884 .await?
885 .context(UnexpectedSnafu {
886 reason: format!("Sink table metadata is missing: {sink_table_name:?}"),
887 })?;
888 let (auto_columns, plan) = match resolve_sink_layout(&inferred_schema, &sink_schema) {
889 Ok(auto_columns) => (auto_columns, flow_plan),
890 Err(normal_error) => {
891 let has_distinct = {
894 let mut found = false;
895 flow_plan
896 .apply(|node| {
897 if matches!(node, LogicalPlan::Distinct(Distinct::All(_))) {
898 found = true;
899 }
900 Ok(datafusion_common::tree_node::TreeNodeRecursion::Continue)
901 })
902 .context(DatafusionSnafu {
903 context: "Failed to inspect flow plan",
904 })?;
905 found
906 };
907 if has_distinct
908 || !sink_exists
909 || !is_explicit_source_timestamp_compatibility(
910 &inferred_schema,
911 &lineage,
912 &sink_schema,
913 &source_schema,
914 )
915 {
916 return Err(normal_error);
917 }
918 let source_timestamp_index = source_schema.timestamp_index().unwrap();
919 let source_timestamp_name =
920 &source_schema.column_schemas()[source_timestamp_index].name;
921 let plan = stateless::rewrite_source_timestamp(
922 flow_plan,
923 &TableReference::full(
924 source_table_name[0].clone(),
925 source_table_name[1].clone(),
926 source_table_name[2].clone(),
927 ),
928 source_timestamp_name,
929 )?;
930 let (effective_output, _) = output_column_schemas(&plan, &source_schema)?;
931 ensure!(
932 effective_output.len() == sink_schema.len()
933 && effective_output
934 .iter()
935 .zip(&sink_schema)
936 .all(|(output, sink)| output.data_type == sink.data_type),
937 InvalidQuerySnafu {
938 reason: "Compatibility plan output does not match the full sink schema"
939 }
940 );
941 (vec![], plan)
942 }
943 };
944
945 Ok(StatelessFlow {
946 source_table_id,
947 source_table_name,
948 source_schema_version: source_schema.version(),
949 source_schema,
950 sink_table_name: sink_table_name.clone(),
951 sink_schema,
952 sink_primary_keys,
953 auto_columns,
954 plan,
955 query_ctx,
956 create_args: args.clone(),
957 })
958 }
959
960 pub async fn flush_flow_inner(&self, _flow_id: FlowId) -> Result<usize, Error> {
961 Ok(0)
962 }
963
964 pub(crate) async fn stateless_flow_ids(&self) -> Vec<FlowId> {
965 let slots = self.stateless_flows.read().await;
966 let mut ids = Vec::new();
967 for (id, slot) in slots.iter() {
968 if slot.runtime.read().await.flow.is_some() {
969 ids.push(*id);
970 }
971 }
972 ids
973 }
974
975 pub async fn flow_exist_inner(&self, flow_id: FlowId) -> Result<bool, Error> {
976 let slot = self.stateless_flows.read().await.get(&flow_id).cloned();
977 Ok(match slot {
978 Some(slot) => slot.runtime.read().await.flow.is_some(),
979 None => false,
980 })
981 }
982
983 async fn fetch_table_pk_schema(
984 &self,
985 table_name: &TableName,
986 ) -> Result<Option<(Vec<String>, Option<usize>, Vec<ColumnSchema>)>, Error> {
987 if let Some(table_id) = self
988 .table_info_source
989 .get_opt_table_id_from_name(table_name)
990 .await?
991 {
992 let table_info = self
993 .table_info_source
994 .get_table_info_value(&table_id)
995 .await?
996 .unwrap();
997 let meta = table_info.table_info.meta;
998 let schema = meta.schema.column_schemas().to_vec();
999 let primary_keys = meta
1000 .primary_key_indices
1001 .into_iter()
1002 .map(|i| schema[i].name.clone())
1003 .collect_vec();
1004 Ok(Some((primary_keys, meta.schema.timestamp_index(), schema)))
1005 } else {
1006 Ok(None)
1007 }
1008 }
1009
1010 async fn adjust_auto_created_table_schema(
1011 &self,
1012 schema: &RelationDesc,
1013 ) -> Result<(Vec<String>, Vec<ColumnSchema>, bool), Error> {
1014 let primary_keys = schema
1015 .typ()
1016 .keys
1017 .first()
1018 .map(|key| {
1019 key.column_indices
1020 .iter()
1021 .map(|i| {
1022 schema
1023 .get_name(*i)
1024 .clone()
1025 .unwrap_or_else(|| format!("col_{i}"))
1026 })
1027 .collect_vec()
1028 })
1029 .unwrap_or_default();
1030 let mut columns = relation_desc_to_column_schemas_with_fallback(schema);
1031 columns.push(ColumnSchema::new(
1032 AUTO_CREATED_UPDATE_AT_TS_COL,
1033 ConcreteDataType::timestamp_millisecond_datatype(),
1034 true,
1035 ));
1036 let no_time_index = schema.typ().time_index.is_none();
1037 if no_time_index {
1038 columns.push(
1039 ColumnSchema::new(
1040 AUTO_CREATED_PLACEHOLDER_TS_COL,
1041 ConcreteDataType::timestamp_millisecond_datatype(),
1042 true,
1043 )
1044 .with_time_index(true),
1045 );
1046 }
1047 Ok((primary_keys, columns, no_time_index))
1048 }
1049}
1050
1051impl StreamingEngine {
1052 async fn handle_inserts_inner(
1053 &self,
1054 request: api::v1::region::InsertRequests,
1055 ) -> Result<(), Error> {
1056 let mut grouped: BTreeMap<_, (_, Vec<DiffRow>, Vec<ConcreteDataType>, u32)> =
1059 BTreeMap::new();
1060 let mut poisoned_tables = std::collections::HashSet::new();
1061 let mut first_error = None;
1062 for write_request in request.requests {
1063 let region_id = RegionId::from(write_request.region_id);
1064 let table_id = region_id.table_id();
1065 if poisoned_tables.contains(&table_id) {
1066 continue;
1067 }
1068 match self.handle_insert_request(write_request).await {
1069 Ok((rows, types, version)) => {
1070 let entry = grouped
1071 .entry(table_id)
1072 .or_insert_with(|| (region_id, vec![], types.clone(), version));
1073 if entry.2 != types || entry.3 != version {
1074 poisoned_tables.insert(table_id);
1075 grouped.remove(&table_id);
1076 let err = InvalidQuerySnafu { reason: format!("Source table {table_id} metadata changed within one insert envelope") }.build();
1077 let ids: Vec<FlowId> = self
1078 .flow_ids_for_table(table_id)
1079 .await
1080 .into_iter()
1081 .map(|(id, _)| id)
1082 .collect();
1083 let err = InsertIntoFlowSnafu {
1084 region_id: u64::from(region_id),
1085 flow_ids: ids,
1086 }
1087 .into_error(BoxedError::new(err));
1088 error!(err; "Failed to normalize flow insert request for region_id={region_id}");
1089 if first_error.is_none() {
1090 first_error = Some(err);
1091 }
1092 } else {
1093 entry.1.extend(rows);
1094 }
1095 }
1096 Err(err) => {
1097 let ids: Vec<FlowId> = self
1098 .flow_ids_for_table(table_id)
1099 .await
1100 .into_iter()
1101 .map(|(id, _)| id)
1102 .collect();
1103 poisoned_tables.insert(table_id);
1104 grouped.remove(&table_id);
1105 let err = InsertIntoFlowSnafu {
1106 region_id: u64::from(region_id),
1107 flow_ids: ids,
1108 }
1109 .into_error(BoxedError::new(err));
1110 error!(err; "Failed to normalize flow insert request for region_id={region_id}");
1111 if first_error.is_none() {
1112 first_error = Some(err);
1113 }
1114 }
1115 }
1116 }
1117 for (_, (region_id, rows, types, version)) in grouped {
1118 if let Err(err) = self
1119 .handle_write_request(region_id, rows, &types, version)
1120 .await
1121 && first_error.is_none()
1122 {
1123 first_error = Some(err);
1124 }
1125 }
1126 match first_error {
1127 Some(err) => Err(err),
1128 None => Ok(()),
1129 }
1130 }
1131
1132 async fn handle_insert_request(
1133 &self,
1134 write_request: api::v1::region::InsertRequest,
1135 ) -> Result<(Vec<DiffRow>, Vec<ConcreteDataType>, u32), Error> {
1136 let region_id = write_request.region_id;
1137 let table_id = RegionId::from(region_id).table_id();
1138 let (insert_schema, rows_proto) = write_request
1139 .rows
1140 .map(|r| (r.schema, r.rows))
1141 .unwrap_or_default();
1142 let now = common_time::util::current_time_millis();
1143 let (table_types, fetch_order, source_schema_version) = {
1144 let table_info = self
1148 .table_info_source
1149 .get_table_info_value(&table_id)
1150 .await?
1151 .context(UnexpectedSnafu {
1152 reason: format!("Table metadata is missing for table id {table_id}"),
1153 })?;
1154 let source_schema_version = table_info.table_info.meta.schema.version();
1155 let table_schema = table_info_value_to_relation_desc(table_info)?;
1156 let defaults = table_schema
1157 .default_values
1158 .iter()
1159 .zip(table_schema.relation_desc.typ().column_types.iter())
1160 .map(|(value, ty)| {
1161 value.as_ref().and_then(|value| {
1162 value.create_default(ty.scalar_type(), ty.nullable()).ok()
1163 })
1164 })
1165 .collect_vec();
1166 let types = table_schema
1167 .relation_desc
1168 .typ()
1169 .column_types
1170 .iter()
1171 .map(|ty| ty.scalar_type.clone())
1172 .collect_vec();
1173 let names = table_schema
1174 .relation_desc
1175 .names
1176 .iter()
1177 .enumerate()
1178 .map(|(idx, name)| {
1179 name.clone().context(InternalSnafu {
1180 reason: format!("Column {idx} of table {table_id} has no name"),
1181 })
1182 })
1183 .collect::<Result<Vec<_>, _>>()?;
1184 let input_columns = insert_schema
1185 .iter()
1186 .enumerate()
1187 .map(|(idx, column)| (&column.column_name, idx))
1188 .collect::<std::collections::HashMap<_, _>>();
1189 let order = names
1190 .iter()
1191 .zip(defaults)
1192 .map(|(name, default)| {
1193 input_columns
1194 .get(name)
1195 .copied()
1196 .map(FetchFromRow::Idx)
1197 .or_else(|| default.map(FetchFromRow::Default))
1198 .with_context(|| UnexpectedSnafu {
1199 reason: format!("Column not found: {name}"),
1200 })
1201 })
1202 .collect::<Result<Vec<_>, _>>()?;
1203 (types, order, source_schema_version)
1204 };
1205 let rows = rows_proto
1206 .into_iter()
1207 .map(|row| {
1208 let row = Row::from(row);
1209 let row = fetch_order
1210 .iter()
1211 .map(|item| item.fetch(&row))
1212 .collect_vec();
1213 (Row::new(row), now, 1)
1214 })
1215 .collect_vec();
1216 Ok((rows, table_types, source_schema_version))
1217 }
1218}
1219
1220#[derive(Debug, Clone)]
1221enum FetchFromRow {
1222 Idx(usize),
1223 Default(datatypes::value::Value),
1224}
1225
1226impl FetchFromRow {
1227 fn fetch(&self, row: &Row) -> datatypes::value::Value {
1228 match self {
1229 Self::Idx(idx) => row
1230 .get(*idx)
1231 .cloned()
1232 .unwrap_or(datatypes::value::Value::Null),
1233 Self::Default(value) => value.clone(),
1234 }
1235 }
1236}