1use std::collections::HashSet;
16use std::ops::Bound;
17
18use bytes::Bytes;
19use futures::future::Either;
20use futures::stream::{self, select_with_strategy};
21use futures_async_stream::try_stream;
22use itertools::Itertools;
23use risingwave_common::array::Op;
24use risingwave_common::bitmap::Bitmap;
25use risingwave_common::catalog::{
26 ColumnDesc, ConflictBehavior, TableId, checked_conflict_behaviors,
27};
28use risingwave_common::hash::{VirtualNode, VnodeBitmapExt};
29use risingwave_common::row::{OwnedRow, RowExt};
30use risingwave_common::types::DataType;
31use risingwave_common::util::sort_util::ColumnOrder;
32use risingwave_common::util::value_encoding::BasicSerde;
33use risingwave_hummock_sdk::HummockReadEpoch;
34use risingwave_pb::catalog::Table;
35use risingwave_pb::catalog::table::Engine;
36use risingwave_pb::id::{SourceId, SubscriberId};
37use risingwave_pb::stream_plan::SubscriptionUpstreamInfo;
38use risingwave_storage::row_serde::value_serde::{ValueRowSerde, ValueRowSerdeNew};
39use risingwave_storage::store::{PrefetchOptions, TryWaitEpochOptions};
40use risingwave_storage::table::KeyedRow;
41
42use crate::common::change_buffer::output_kind as cb_kind;
43use crate::common::metrics::MetricsInfo;
44use crate::common::table::state_table::{
45 StateTableBuilder, StateTableInner, StateTableOpConsistencyLevel,
46};
47use crate::executor::error::ErrorKind;
48use crate::executor::monitor::MaterializeMetrics;
49use crate::executor::mview::RefreshProgressTable;
50use crate::executor::mview::cache::MaterializeCache;
51use crate::executor::prelude::*;
52use crate::executor::{BarrierInner, BarrierMutationType, EpochPair};
53use crate::task::LocalBarrierManager;
54
55#[derive(Debug, Clone)]
56pub enum MaterializeStreamState<M> {
57 NormalIngestion,
58 MergingData,
59 CleanUp,
60 CommitAndYieldBarrier {
61 barrier: BarrierInner<M>,
62 expect_next_state: Box<MaterializeStreamState<M>>,
63 },
64 RefreshEnd {
65 on_complete_epoch: EpochPair,
66 },
67}
68
69pub struct MaterializeExecutor<S: StateStore, SD: ValueRowSerde> {
71 input: Executor,
72
73 schema: Schema,
74
75 state_table: StateTableInner<S, SD>,
76
77 stream_key_indices: Vec<usize>,
88
89 arrange_key_indices: Vec<usize>,
91
92 actor_context: ActorContextRef,
93
94 materialize_cache: Option<MaterializeCache>,
96
97 toastable_column_indices: Option<Vec<usize>>,
102
103 conflict_behavior: ConflictBehavior,
104
105 cleaned_by_ttl_watermark: bool,
107
108 version_column_indices: Vec<u32>,
109
110 may_have_downstream: bool,
111
112 subscriber_ids: HashSet<SubscriberId>,
113
114 metrics: MaterializeMetrics,
115
116 is_dummy_table: bool,
119
120 refresh_args: Option<RefreshableMaterializeArgs<S, SD>>,
122
123 local_barrier_manager: LocalBarrierManager,
125}
126
127pub struct RefreshableMaterializeArgs<S: StateStore, SD: ValueRowSerde> {
129 pub table_catalog: Table,
131
132 pub staging_table_catalog: Table,
134
135 pub is_refreshing: bool,
137
138 pub staging_table: StateTableInner<S, SD>,
145
146 pub progress_table: RefreshProgressTable<S>,
148
149 pub table_id: TableId,
151}
152
153impl<S: StateStore, SD: ValueRowSerde> RefreshableMaterializeArgs<S, SD> {
154 pub async fn new(
156 store: S,
157 table_catalog: &Table,
158 staging_table_catalog: &Table,
159 progress_state_table: &Table,
160 vnodes: Option<Arc<Bitmap>>,
161 ) -> Self {
162 let table_id = table_catalog.id;
163
164 let staging_table = StateTableInner::from_table_catalog_inconsistent_op(
166 staging_table_catalog,
167 store.clone(),
168 vnodes.clone(),
169 )
170 .await;
171
172 let progress_state_table = StateTableInner::from_table_catalog_inconsistent_op(
173 progress_state_table,
174 store,
175 vnodes,
176 )
177 .await;
178
179 let pk_len = table_catalog.pk.len();
181 let progress_table = RefreshProgressTable::new(progress_state_table, pk_len);
182
183 debug_assert_eq!(staging_table.vnodes(), progress_table.vnodes());
184
185 Self {
186 table_catalog: table_catalog.clone(),
187 staging_table_catalog: staging_table_catalog.clone(),
188 is_refreshing: false,
189 staging_table,
190 progress_table,
191 table_id,
192 }
193 }
194}
195
196fn get_op_consistency_level(
197 cleaned_by_ttl_watermark: bool,
198 conflict_behavior: ConflictBehavior,
199 may_have_downstream: bool,
200 subscriber_ids: &HashSet<SubscriberId>,
201) -> StateTableOpConsistencyLevel {
202 if cleaned_by_ttl_watermark {
203 StateTableOpConsistencyLevel::Inconsistent
207 } else if !subscriber_ids.is_empty() {
208 StateTableOpConsistencyLevel::LogStoreEnabled
209 } else if !may_have_downstream && matches!(conflict_behavior, ConflictBehavior::Overwrite) {
210 StateTableOpConsistencyLevel::Inconsistent
213 } else {
214 StateTableOpConsistencyLevel::ConsistentOldValue
215 }
216}
217
218impl<S: StateStore, SD: ValueRowSerde> MaterializeExecutor<S, SD> {
219 #[expect(clippy::too_many_arguments)]
223 pub async fn new(
224 input: Executor,
225 schema: Schema,
226 store: S,
227 arrange_key: Vec<ColumnOrder>,
228 actor_context: ActorContextRef,
229 vnodes: Option<Arc<Bitmap>>,
230 table_catalog: &Table,
231 watermark_epoch: AtomicU64Ref,
232 conflict_behavior: ConflictBehavior,
233 version_column_indices: Vec<u32>,
234 metrics: Arc<StreamingMetrics>,
235 refresh_args: Option<RefreshableMaterializeArgs<S, SD>>,
236 cleaned_by_ttl_watermark: bool,
237 local_barrier_manager: LocalBarrierManager,
238 ) -> Self {
239 let table_columns: Vec<ColumnDesc> = table_catalog
240 .columns
241 .iter()
242 .map(|col| col.column_desc.as_ref().unwrap().into())
243 .collect();
244
245 let toastable_column_indices = if table_catalog.cdc_table_type()
248 == risingwave_pb::catalog::table::CdcTableType::Postgres
249 {
250 let toastable_indices: Vec<usize> = table_columns
251 .iter()
252 .enumerate()
253 .filter_map(|(index, column)| match &column.data_type {
254 DataType::Varchar
264 | DataType::List(_)
265 | DataType::Bytea
266 | DataType::Jsonb
267 | DataType::Vector(_) => Some(index),
268 _ => None,
269 })
270 .collect();
271
272 if toastable_indices.is_empty() {
273 None
274 } else {
275 Some(toastable_indices)
276 }
277 } else {
278 None
279 };
280
281 let row_serde: BasicSerde = BasicSerde::new(
282 Arc::from_iter(table_catalog.value_indices.iter().map(|val| *val as usize)),
283 Arc::from(table_columns.into_boxed_slice()),
284 );
285
286 let stream_key_indices: Vec<usize> = table_catalog
287 .stream_key
288 .iter()
289 .map(|idx| *idx as usize)
290 .collect();
291 let arrange_key_indices: Vec<usize> = arrange_key.iter().map(|k| k.column_index).collect();
292 let may_have_downstream = actor_context.initial_dispatch_num != 0;
293 let subscriber_ids = actor_context.initial_subscriber_ids.clone();
294 let op_consistency_level = get_op_consistency_level(
295 cleaned_by_ttl_watermark,
296 conflict_behavior,
297 may_have_downstream,
298 &subscriber_ids,
299 );
300 let state_table_metrics = metrics.new_state_table_metrics(
301 table_catalog.id,
302 actor_context.id,
303 actor_context.fragment_id,
304 );
305 let state_table = StateTableBuilder::new(table_catalog, store, vnodes)
307 .with_op_consistency_level(op_consistency_level)
308 .enable_preload_all_rows_by_config(&actor_context.config)
309 .enable_vnode_key_stats(
310 actor_context
311 .config
312 .developer
313 .enable_vnode_key_stats_for_materialize,
314 &actor_context.config,
315 )
316 .with_metrics(state_table_metrics)
317 .build()
318 .await;
319
320 let mv_metrics = metrics.new_materialize_metrics(
321 table_catalog.id,
322 actor_context.id,
323 actor_context.fragment_id,
324 );
325 let cache_metrics = metrics.new_materialize_cache_metrics(
326 table_catalog.id,
327 actor_context.id,
328 actor_context.fragment_id,
329 );
330
331 let metrics_info =
332 MetricsInfo::new(metrics, table_catalog.id, actor_context.id, "Materialize");
333
334 let is_dummy_table =
335 table_catalog.engine == Some(Engine::Iceberg as i32) && table_catalog.append_only;
336
337 Self {
338 input,
339 schema,
340 state_table,
341 stream_key_indices,
342 arrange_key_indices,
343 actor_context,
344 materialize_cache: MaterializeCache::new(
345 watermark_epoch,
346 metrics_info,
347 row_serde,
348 version_column_indices.clone(),
349 conflict_behavior,
350 toastable_column_indices.clone(),
351 cache_metrics,
352 ),
353 toastable_column_indices,
354 conflict_behavior,
355 cleaned_by_ttl_watermark,
356 version_column_indices,
357 is_dummy_table,
358 may_have_downstream,
359 subscriber_ids,
360 metrics: mv_metrics,
361 refresh_args,
362 local_barrier_manager,
363 }
364 }
365
366 #[try_stream(ok = Message, error = StreamExecutorError)]
367 async fn execute_inner(mut self) {
368 let mv_table_id = self.state_table.table_id();
369 let data_types = self.schema.data_types();
370 let mut input = self.input.execute();
371
372 let barrier = expect_first_barrier(&mut input).await?;
373 let first_epoch = barrier.epoch;
374 let _barrier_epoch = barrier.epoch; yield Message::Barrier(barrier);
377 self.state_table.init_epoch(first_epoch).await?;
378
379 let mut inner_state =
381 Box::new(MaterializeStreamState::<BarrierMutationType>::NormalIngestion);
382 if let Some(ref mut refresh_args) = self.refresh_args {
384 refresh_args.staging_table.init_epoch(first_epoch).await?;
385
386 refresh_args.progress_table.recover(first_epoch).await?;
388
389 let progress_stats = refresh_args.progress_table.get_progress_stats();
391 if progress_stats.total_vnodes > 0 && !progress_stats.is_complete() {
392 refresh_args.is_refreshing = true;
393 tracing::info!(
394 total_vnodes = progress_stats.total_vnodes,
395 completed_vnodes = progress_stats.completed_vnodes,
396 "Recovered refresh in progress, resuming refresh operation"
397 );
398
399 let incomplete_vnodes: Vec<_> = refresh_args
403 .progress_table
404 .get_all_progress()
405 .iter()
406 .filter(|(_, entry)| !entry.is_completed)
407 .map(|(&vnode, _)| vnode)
408 .collect();
409
410 if !incomplete_vnodes.is_empty() {
411 tracing::info!(
413 incomplete_vnodes = incomplete_vnodes.len(),
414 "Recovery detected incomplete VNodes, resuming refresh operation"
415 );
416 } else {
419 tracing::warn!("Unexpected recovery state: no incomplete VNodes found");
421 }
422 }
423 }
424
425 if let Some(ref refresh_args) = self.refresh_args
427 && refresh_args.is_refreshing
428 {
429 let incomplete_vnodes: Vec<_> = refresh_args
431 .progress_table
432 .get_all_progress()
433 .iter()
434 .filter(|(_, entry)| !entry.is_completed)
435 .map(|(&vnode, _)| vnode)
436 .collect();
437 if !incomplete_vnodes.is_empty() {
438 *inner_state = MaterializeStreamState::<_>::MergingData;
440 tracing::info!(
441 incomplete_vnodes = incomplete_vnodes.len(),
442 "Recovery: Resuming refresh from merge stage due to incomplete VNodes"
443 );
444 }
445 }
446
447 'main_loop: loop {
449 match *inner_state {
450 MaterializeStreamState::NormalIngestion => {
451 #[for_await]
452 '_normal_ingest: for msg in input.by_ref() {
453 let msg = msg?;
454 if let Some(cache) = &mut self.materialize_cache {
455 cache.evict();
456 }
457
458 match msg {
459 Message::Watermark(w) => {
460 if self.cleaned_by_ttl_watermark
461 && self.state_table.clean_watermark_index == Some(w.col_idx)
462 {
463 self.state_table.update_watermark(w.val.clone());
464 }
465 yield Message::Watermark(w);
466 }
467 Message::Chunk(chunk) if self.is_dummy_table => {
468 self.metrics
469 .materialize_input_row_count
470 .inc_by(chunk.cardinality() as u64);
471 yield Message::Chunk(chunk);
472 }
473 Message::Chunk(chunk) => {
474 self.metrics
475 .materialize_input_row_count
476 .inc_by(chunk.cardinality() as u64);
477
478 let optimized_conflict_behavior = if let ConflictBehavior::Overwrite =
491 self.conflict_behavior
492 && !self.state_table.is_consistent_op()
493 && !self.cleaned_by_ttl_watermark
494 && self.version_column_indices.is_empty()
495 && self.toastable_column_indices.is_none()
496 {
497 ConflictBehavior::NoCheck
498 } else {
499 self.conflict_behavior
500 };
501
502 match optimized_conflict_behavior {
503 checked_conflict_behaviors!() => {
504 if chunk.cardinality() == 0 {
505 continue;
507 }
508
509 if let Some(ref mut refresh_args) = self.refresh_args
512 && refresh_args.is_refreshing
513 {
514 let key_chunk = chunk
515 .clone()
516 .project(self.state_table.pk_indices());
517 tracing::trace!(
518 staging_chunk = %key_chunk.to_pretty(),
519 input_chunk = %chunk.to_pretty(),
520 "writing to staging table"
521 );
522 if cfg!(debug_assertions) {
523 assert!(
525 key_chunk
526 .ops()
527 .iter()
528 .all(|op| op == &Op::Insert)
529 );
530 }
531 refresh_args
532 .staging_table
533 .write_chunk(key_chunk.clone());
534 refresh_args.staging_table.try_flush().await?;
535 }
536
537 let cache = self.materialize_cache.as_mut().unwrap();
538 let change_buffer =
539 cache.handle_new(chunk, &self.state_table).await?;
540
541 let output_chunk = if self.stream_key_indices
542 == self.state_table.pk_indices()
543 {
544 change_buffer.into_chunk::<{ cb_kind::RETRACT }>(
545 data_types.clone(),
546 )
547 } else {
548 debug_assert!(
551 self.state_table
552 .pk_indices()
553 .iter()
554 .all(|&i| self.stream_key_indices.contains(&i))
555 );
556 change_buffer.into_chunk_with_key(
557 data_types.clone(),
558 &self.stream_key_indices,
559 )
560 };
561
562 match output_chunk {
563 Some(output_chunk) => {
564 self.state_table.write_chunk(output_chunk.clone());
565 self.state_table.try_flush().await?;
566 yield Message::Chunk(output_chunk);
567 }
568 None => continue,
569 }
570 }
571 ConflictBehavior::NoCheck => {
572 self.state_table.write_chunk(chunk.clone());
573 self.state_table.try_flush().await?;
574
575 if let Some(ref mut refresh_args) = self.refresh_args
577 && refresh_args.is_refreshing
578 {
579 let key_chunk = chunk
580 .clone()
581 .project(self.state_table.pk_indices());
582 tracing::trace!(
583 staging_chunk = %key_chunk.to_pretty(),
584 input_chunk = %chunk.to_pretty(),
585 "writing to staging table"
586 );
587 if cfg!(debug_assertions) {
588 assert!(
590 key_chunk
591 .ops()
592 .iter()
593 .all(|op| op == &Op::Insert)
594 );
595 }
596 refresh_args
597 .staging_table
598 .write_chunk(key_chunk.clone());
599 refresh_args.staging_table.try_flush().await?;
600 }
601
602 yield Message::Chunk(chunk);
603 }
604 }
605 }
606 Message::Barrier(barrier) => {
607 *inner_state = MaterializeStreamState::CommitAndYieldBarrier {
608 barrier,
609 expect_next_state: Box::new(
610 MaterializeStreamState::NormalIngestion,
611 ),
612 };
613 continue 'main_loop;
614 }
615 }
616 }
617
618 return Err(StreamExecutorError::from(ErrorKind::Uncategorized(
619 anyhow::anyhow!(
620 "Input stream terminated unexpectedly during normal ingestion"
621 ),
622 )));
623 }
624 MaterializeStreamState::MergingData => {
625 let Some(refresh_args) = self.refresh_args.as_mut() else {
626 panic!(
627 "MaterializeExecutor entered CleanUp state without refresh_args configured"
628 );
629 };
630 tracing::info!(table_id = %refresh_args.table_id, "on_load_finish: Starting table replacement operation");
631
632 debug_assert_eq!(
633 self.state_table.vnodes(),
634 refresh_args.staging_table.vnodes()
635 );
636 debug_assert_eq!(
637 refresh_args.staging_table.vnodes(),
638 refresh_args.progress_table.vnodes()
639 );
640
641 let mut rows_to_delete = vec![];
642 let mut merge_complete = false;
643 let mut pending_barrier: Option<Barrier> = None;
644
645 {
647 let left_input = input.by_ref().map(Either::Left);
648 let right_merge_sort = pin!(
649 Self::make_mergesort_stream(
650 &self.state_table,
651 &refresh_args.staging_table,
652 &mut refresh_args.progress_table
653 )
654 .map(Either::Right)
655 );
656
657 let mut merge_stream =
660 select_with_strategy(left_input, right_merge_sort, |_: &mut ()| {
661 stream::PollNext::Left
662 });
663
664 #[for_await]
665 'merge_stream: for either in &mut merge_stream {
666 match either {
667 Either::Left(msg) => {
668 let msg = msg?;
669 match msg {
670 Message::Watermark(w) => yield Message::Watermark(w),
671 Message::Chunk(chunk) => {
672 tracing::warn!(chunk = %chunk.to_pretty(), "chunk is ignored during merge phase");
673 }
674 Message::Barrier(b) => {
675 pending_barrier = Some(b);
676 break 'merge_stream;
677 }
678 }
679 }
680 Either::Right(result) => {
681 match result? {
682 Some((_vnode, row)) => {
683 rows_to_delete.push(row);
684 }
685 None => {
686 merge_complete = true;
688
689 }
691 }
692 }
693 }
694 }
695 }
696
697 for row in &rows_to_delete {
699 self.state_table.delete(row);
700 }
701 if let Some(cache) = &mut self.materialize_cache {
702 cache.invalidate_rows(&rows_to_delete, &self.state_table);
703 }
704 if !rows_to_delete.is_empty() {
705 let to_delete_chunk = StreamChunk::from_rows(
706 &rows_to_delete
707 .iter()
708 .map(|row| (Op::Delete, row))
709 .collect_vec(),
710 &self.schema.data_types(),
711 );
712
713 yield Message::Chunk(to_delete_chunk);
714 }
715
716 assert!(pending_barrier.is_some(), "pending barrier is not set");
718
719 *inner_state = MaterializeStreamState::CommitAndYieldBarrier {
720 barrier: pending_barrier.unwrap(),
721 expect_next_state: if merge_complete {
722 Box::new(MaterializeStreamState::CleanUp)
723 } else {
724 Box::new(MaterializeStreamState::MergingData)
725 },
726 };
727 continue 'main_loop;
728 }
729 MaterializeStreamState::CleanUp => {
730 let Some(refresh_args) = self.refresh_args.as_mut() else {
731 panic!(
732 "MaterializeExecutor entered MergingData state without refresh_args configured"
733 );
734 };
735 tracing::info!(table_id = %refresh_args.table_id, "on_load_finish: resuming CleanUp Stage");
736
737 #[for_await]
738 for msg in input.by_ref() {
739 let msg = msg?;
740 match msg {
741 Message::Watermark(w) => {
742 if self.cleaned_by_ttl_watermark
743 && self.state_table.clean_watermark_index == Some(w.col_idx)
744 {
745 self.state_table.update_watermark(w.val.clone());
746 }
747 yield Message::Watermark(w)
748 }
749 Message::Chunk(chunk) => {
750 tracing::warn!(chunk = %chunk.to_pretty(), "chunk is ignored during merge phase");
751 }
752 Message::Barrier(barrier) if !barrier.is_checkpoint() => {
753 *inner_state = MaterializeStreamState::CommitAndYieldBarrier {
754 barrier,
755 expect_next_state: Box::new(MaterializeStreamState::CleanUp),
756 };
757 continue 'main_loop;
758 }
759 Message::Barrier(barrier) => {
760 let staging_table_id = refresh_args.staging_table.table_id();
761 let epoch = barrier.epoch;
762 self.local_barrier_manager.report_refresh_finished(
763 epoch,
764 self.actor_context.id,
765 refresh_args.table_id,
766 staging_table_id,
767 );
768 tracing::debug!(table_id = %refresh_args.table_id, "on_load_finish: Reported staging table truncation and diff applied");
769
770 *inner_state = MaterializeStreamState::CommitAndYieldBarrier {
771 barrier,
772 expect_next_state: Box::new(
773 MaterializeStreamState::RefreshEnd {
774 on_complete_epoch: epoch,
775 },
776 ),
777 };
778 continue 'main_loop;
779 }
780 }
781 }
782 }
783 MaterializeStreamState::RefreshEnd { on_complete_epoch } => {
784 let Some(refresh_args) = self.refresh_args.as_mut() else {
785 panic!(
786 "MaterializeExecutor entered RefreshEnd state without refresh_args configured"
787 );
788 };
789 let staging_table_id = refresh_args.staging_table.table_id();
790
791 let staging_store = refresh_args.staging_table.state_store().clone();
793 staging_store
794 .try_wait_epoch(
795 HummockReadEpoch::Committed(on_complete_epoch.prev),
796 TryWaitEpochOptions {
797 table_id: staging_table_id,
798 },
799 )
800 .await?;
801
802 tracing::info!(table_id = %refresh_args.table_id, "RefreshEnd: Refresh completed");
803
804 if let Some(ref mut refresh_args) = self.refresh_args {
805 refresh_args.is_refreshing = false;
806 }
807 *inner_state = MaterializeStreamState::NormalIngestion;
808 continue 'main_loop;
809 }
810 MaterializeStreamState::CommitAndYieldBarrier {
811 barrier,
812 mut expect_next_state,
813 } => {
814 if let Some(ref mut refresh_args) = self.refresh_args {
815 match barrier.mutation.as_deref() {
816 Some(Mutation::RefreshStart {
817 table_id: refresh_table_id,
818 associated_source_id: _,
819 }) if *refresh_table_id == refresh_args.table_id => {
820 debug_assert!(
821 !refresh_args.is_refreshing,
822 "cannot start refresh twice"
823 );
824 refresh_args.is_refreshing = true;
825 tracing::info!(table_id = %refresh_table_id, "RefreshStart barrier received");
826
827 Self::init_refresh_progress(
829 &self.state_table,
830 &mut refresh_args.progress_table,
831 barrier.epoch.curr,
832 )?;
833 }
834 Some(Mutation::LoadFinish {
835 associated_source_id: load_finish_source_id,
836 }) => {
837 let associated_source_id: SourceId = match refresh_args
839 .table_catalog
840 .optional_associated_source_id
841 {
842 Some(id) => id.into(),
843 None => unreachable!("associated_source_id is not set"),
844 };
845
846 if *load_finish_source_id == associated_source_id {
847 tracing::info!(
848 %load_finish_source_id,
849 "LoadFinish received, starting data replacement"
850 );
851 *expect_next_state = MaterializeStreamState::<_>::MergingData;
852 }
853 }
854 _ => {}
855 }
856 }
857
858 if !self.may_have_downstream
862 && barrier.has_more_downstream_fragments(self.actor_context.id)
863 {
864 self.may_have_downstream = true;
865 }
866 Self::may_update_depended_subscriptions(
867 &mut self.subscriber_ids,
868 &barrier,
869 mv_table_id,
870 );
871 let op_consistency_level = get_op_consistency_level(
872 self.cleaned_by_ttl_watermark,
873 self.conflict_behavior,
874 self.may_have_downstream,
875 &self.subscriber_ids,
876 );
877 let post_commit = self
878 .state_table
879 .commit_may_switch_consistent_op(barrier.epoch, op_consistency_level)
880 .await?;
881
882 let update_vnode_bitmap = barrier.as_update_vnode_bitmap(self.actor_context.id);
883
884 let refresh_post_commit = if let Some(ref mut refresh_args) = self.refresh_args
886 {
887 Some((
890 refresh_args.staging_table.commit(barrier.epoch).await?,
891 refresh_args.progress_table.commit(barrier.epoch).await?,
892 ))
893 } else {
894 None
895 };
896
897 let b_epoch = barrier.epoch;
898 yield Message::Barrier(barrier);
899
900 if let Some((_, cache_may_stale)) = post_commit
902 .post_yield_barrier(update_vnode_bitmap.clone())
903 .await?
904 && cache_may_stale
905 && let Some(cache) = &mut self.materialize_cache
906 {
907 cache.clear();
908 }
909
910 if let Some((staging_post_commit, progress_post_commit)) = refresh_post_commit {
912 staging_post_commit
913 .post_yield_barrier(update_vnode_bitmap.clone())
914 .await?;
915 progress_post_commit
916 .post_yield_barrier(update_vnode_bitmap)
917 .await?;
918 }
919
920 self.metrics
921 .materialize_current_epoch
922 .set(b_epoch.curr as i64);
923
924 *inner_state = *expect_next_state;
927 }
928 }
929 }
930 }
931
932 #[try_stream(ok = Option<(VirtualNode, OwnedRow)>, error = StreamExecutorError)]
936 async fn make_mergesort_stream<'a>(
937 main_table: &'a StateTableInner<S, SD>,
938 staging_table: &'a StateTableInner<S, SD>,
939 progress_table: &'a mut RefreshProgressTable<S>,
940 ) {
941 for vnode in main_table.vnodes().clone().iter_vnodes() {
942 let mut processed_rows = 0;
943 let pk_range: (Bound<OwnedRow>, Bound<OwnedRow>) =
945 if let Some(current_entry) = progress_table.get_progress(vnode) {
946 if current_entry.is_completed {
948 tracing::debug!(
949 vnode = vnode.to_index(),
950 "Skipping already completed VNode during recovery"
951 );
952 continue;
953 }
954 processed_rows += current_entry.processed_rows;
955 tracing::debug!(vnode = vnode.to_index(), "Started merging VNode");
956
957 if let Some(current_state) = ¤t_entry.current_pos {
958 (Bound::Excluded(current_state.clone()), Bound::Unbounded)
959 } else {
960 (Bound::Unbounded, Bound::Unbounded)
961 }
962 } else {
963 (Bound::Unbounded, Bound::Unbounded)
964 };
965
966 let iter_main = main_table
967 .iter_keyed_row_with_vnode(
968 vnode,
969 &pk_range,
970 PrefetchOptions::prefetch_for_large_range_scan(),
971 )
972 .await?;
973 let iter_staging = staging_table
974 .iter_keyed_row_with_vnode(
975 vnode,
976 &pk_range,
977 PrefetchOptions::prefetch_for_large_range_scan(),
978 )
979 .await?;
980
981 pin_mut!(iter_main);
982 pin_mut!(iter_staging);
983
984 let mut main_item: Option<KeyedRow<Bytes>> = iter_main.next().await.transpose()?;
986 let mut staging_item: Option<KeyedRow<Bytes>> =
987 iter_staging.next().await.transpose()?;
988
989 while let Some(main_kv) = main_item {
990 let main_key = main_kv.key();
991
992 let mut should_delete = false;
994 while let Some(staging_kv) = &staging_item {
995 let staging_key = staging_kv.key();
996 match main_key.cmp(staging_key) {
997 std::cmp::Ordering::Greater => {
998 staging_item = iter_staging.next().await.transpose()?;
1000 }
1001 std::cmp::Ordering::Equal => {
1002 break;
1004 }
1005 std::cmp::Ordering::Less => {
1006 should_delete = true;
1008 break;
1009 }
1010 }
1011 }
1012
1013 if staging_item.is_none() {
1015 should_delete = true;
1016 }
1017
1018 if should_delete {
1019 yield Some((vnode, main_kv.row().clone()));
1020 }
1021
1022 processed_rows += 1;
1024 tracing::debug!(
1025 "set progress table: vnode = {:?}, processed_rows = {:?}",
1026 vnode,
1027 processed_rows
1028 );
1029 progress_table.set_progress(
1030 vnode,
1031 Some(
1032 main_kv
1033 .row()
1034 .project(main_table.pk_indices())
1035 .to_owned_row(),
1036 ),
1037 false,
1038 processed_rows,
1039 )?;
1040 main_item = iter_main.next().await.transpose()?;
1041 }
1042
1043 if let Some(current_entry) = progress_table.get_progress(vnode) {
1045 progress_table.set_progress(
1046 vnode,
1047 current_entry.current_pos.clone(),
1048 true, current_entry.processed_rows,
1050 )?;
1051
1052 tracing::debug!(vnode = vnode.to_index(), "Completed merging VNode");
1053 }
1054 }
1055
1056 yield None;
1058 }
1059
1060 fn may_update_depended_subscriptions(
1062 depended_subscriptions: &mut HashSet<SubscriberId>,
1063 barrier: &Barrier,
1064 mv_table_id: TableId,
1065 ) {
1066 for subscriber_id in barrier.added_subscriber_on_mv_table(mv_table_id) {
1067 if !depended_subscriptions.insert(subscriber_id) {
1068 warn!(
1069 ?depended_subscriptions,
1070 %mv_table_id,
1071 %subscriber_id,
1072 "subscription id already exists"
1073 );
1074 }
1075 }
1076
1077 if let Some(subscriptions_to_drop) = barrier.as_subscriptions_to_drop() {
1078 for SubscriptionUpstreamInfo {
1079 subscriber_id,
1080 upstream_mv_table_id,
1081 } in subscriptions_to_drop
1082 {
1083 if *upstream_mv_table_id == mv_table_id
1084 && !depended_subscriptions.remove(subscriber_id)
1085 {
1086 warn!(
1087 ?depended_subscriptions,
1088 %mv_table_id,
1089 %subscriber_id,
1090 "drop non existing subscriber_id id"
1091 );
1092 }
1093 }
1094 }
1095 }
1096
1097 fn init_refresh_progress(
1099 state_table: &StateTableInner<S, SD>,
1100 progress_table: &mut RefreshProgressTable<S>,
1101 _epoch: u64,
1102 ) -> StreamExecutorResult<()> {
1103 debug_assert_eq!(state_table.vnodes(), progress_table.vnodes());
1104
1105 for vnode in state_table.vnodes().iter_vnodes() {
1107 progress_table.set_progress(
1108 vnode, None, false, 0, )?;
1112 }
1113
1114 tracing::info!(
1115 vnodes_count = state_table.vnodes().count_ones(),
1116 "Initialized refresh progress tracking for all VNodes"
1117 );
1118
1119 Ok(())
1120 }
1121}
1122
1123impl<S: StateStore> MaterializeExecutor<S, BasicSerde> {
1124 #[cfg(any(test, feature = "test"))]
1126 pub async fn for_test(
1127 input: Executor,
1128 store: S,
1129 table_id: TableId,
1130 keys: Vec<ColumnOrder>,
1131 column_ids: Vec<risingwave_common::catalog::ColumnId>,
1132 watermark_epoch: AtomicU64Ref,
1133 conflict_behavior: ConflictBehavior,
1134 ) -> Self {
1135 Self::for_test_inner(
1136 input,
1137 store,
1138 table_id,
1139 keys,
1140 column_ids,
1141 watermark_epoch,
1142 conflict_behavior,
1143 None,
1144 )
1145 .await
1146 }
1147
1148 #[cfg(any(test, feature = "test"))]
1149 #[expect(clippy::too_many_arguments)]
1150 pub async fn for_test_with_stream_key(
1151 input: Executor,
1152 store: S,
1153 table_id: TableId,
1154 keys: Vec<ColumnOrder>,
1155 stream_key: Vec<usize>,
1156 column_ids: Vec<risingwave_common::catalog::ColumnId>,
1157 watermark_epoch: AtomicU64Ref,
1158 conflict_behavior: ConflictBehavior,
1159 ) -> Self {
1160 Self::for_test_inner(
1161 input,
1162 store,
1163 table_id,
1164 keys,
1165 column_ids,
1166 watermark_epoch,
1167 conflict_behavior,
1168 Some(stream_key),
1169 )
1170 .await
1171 }
1172
1173 #[cfg(any(test, feature = "test"))]
1174 #[expect(clippy::too_many_arguments)]
1175 async fn for_test_inner(
1176 input: Executor,
1177 store: S,
1178 table_id: TableId,
1179 keys: Vec<ColumnOrder>,
1180 column_ids: Vec<risingwave_common::catalog::ColumnId>,
1181 watermark_epoch: AtomicU64Ref,
1182 conflict_behavior: ConflictBehavior,
1183 stream_key: Option<Vec<usize>>,
1184 ) -> Self {
1185 use risingwave_common::util::iter_util::ZipEqFast;
1186
1187 let arrange_columns: Vec<usize> = keys.iter().map(|k| k.column_index).collect();
1188 let arrange_order_types = keys.iter().map(|k| k.order_type).collect();
1189 let schema = input.schema().clone();
1190 let columns: Vec<ColumnDesc> = column_ids
1191 .into_iter()
1192 .zip_eq_fast(schema.fields.iter())
1193 .map(|(column_id, field)| ColumnDesc::unnamed(column_id, field.data_type()))
1194 .collect_vec();
1195
1196 let row_serde = BasicSerde::new(
1197 Arc::from((0..columns.len()).collect_vec()),
1198 Arc::from(columns.clone().into_boxed_slice()),
1199 );
1200 let stream_key_indices = stream_key.unwrap_or_else(|| arrange_columns.clone());
1201 let mut table_catalog = crate::common::table::test_utils::gen_pbtable(
1202 table_id,
1203 columns,
1204 arrange_order_types,
1205 arrange_columns.clone(),
1206 0,
1207 );
1208 table_catalog.stream_key = stream_key_indices.iter().map(|i| *i as i32).collect();
1209 let state_table = StateTableInner::from_table_catalog(&table_catalog, store, None).await;
1210
1211 let unused = StreamingMetrics::unused();
1212 let metrics = unused.new_materialize_metrics(table_id, 1.into(), 2.into());
1213 let cache_metrics = unused.new_materialize_cache_metrics(table_id, 1.into(), 2.into());
1214
1215 Self {
1216 input,
1217 schema,
1218 state_table,
1219 stream_key_indices,
1220 arrange_key_indices: arrange_columns.clone(),
1221 actor_context: ActorContext::for_test(0),
1222 materialize_cache: MaterializeCache::new(
1223 watermark_epoch,
1224 MetricsInfo::for_test(),
1225 row_serde,
1226 vec![],
1227 conflict_behavior,
1228 None,
1229 cache_metrics,
1230 ),
1231 toastable_column_indices: None,
1232 conflict_behavior,
1233 cleaned_by_ttl_watermark: false,
1234 version_column_indices: vec![],
1235 is_dummy_table: false,
1236 may_have_downstream: true,
1237 subscriber_ids: HashSet::new(),
1238 metrics,
1239 refresh_args: None, local_barrier_manager: LocalBarrierManager::for_test(),
1241 }
1242 }
1243}
1244
1245impl<S: StateStore, SD: ValueRowSerde> Execute for MaterializeExecutor<S, SD> {
1246 fn execute(self: Box<Self>) -> BoxedMessageStream {
1247 self.execute_inner().boxed()
1248 }
1249}
1250
1251impl<S: StateStore, SD: ValueRowSerde> std::fmt::Debug for MaterializeExecutor<S, SD> {
1252 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1253 f.debug_struct("MaterializeExecutor")
1254 .field("arrange_key_indices", &self.arrange_key_indices)
1255 .field("stream_key_indices", &self.stream_key_indices)
1256 .finish()
1257 }
1258}
1259
1260#[cfg(test)]
1261mod tests {
1262
1263 use std::iter;
1264 use std::sync::atomic::AtomicU64;
1265
1266 use rand::rngs::SmallRng;
1267 use rand::{Rng, RngCore, SeedableRng};
1268 use risingwave_common::array::stream_chunk::{StreamChunkMut, StreamChunkTestExt};
1269 use risingwave_common::catalog::Field;
1270 use risingwave_common::util::epoch::test_epoch;
1271 use risingwave_common::util::sort_util::OrderType;
1272 use risingwave_hummock_sdk::HummockReadEpoch;
1273 use risingwave_storage::memory::MemoryStateStore;
1274 use risingwave_storage::table::batch_table::BatchTable;
1275
1276 use super::*;
1277 use crate::executor::test_utils::*;
1278
1279 #[tokio::test]
1280 async fn test_materialize_executor() {
1281 let memory_state_store = MemoryStateStore::new();
1283 let table_id = TableId::new(1);
1284 let schema = Schema::new(vec![
1286 Field::unnamed(DataType::Int32),
1287 Field::unnamed(DataType::Int32),
1288 ]);
1289 let column_ids = vec![0.into(), 1.into()];
1290
1291 let chunk1 = StreamChunk::from_pretty(
1293 " i i
1294 + 1 4
1295 + 2 5
1296 + 3 6",
1297 );
1298 let chunk2 = StreamChunk::from_pretty(
1299 " i i
1300 + 7 8
1301 - 3 6",
1302 );
1303
1304 let source = MockSource::with_messages(vec![
1306 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1307 Message::Chunk(chunk1),
1308 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1309 Message::Chunk(chunk2),
1310 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
1311 ])
1312 .into_executor(schema.clone(), StreamKey::new());
1313
1314 let order_types = vec![OrderType::ascending()];
1315 let column_descs = vec![
1316 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
1317 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
1318 ];
1319
1320 let table = BatchTable::for_test(
1321 memory_state_store.clone(),
1322 table_id,
1323 column_descs,
1324 order_types,
1325 vec![0],
1326 vec![0, 1],
1327 );
1328
1329 let mut materialize_executor = MaterializeExecutor::for_test(
1330 source,
1331 memory_state_store,
1332 table_id,
1333 vec![ColumnOrder::new(0, OrderType::ascending())],
1334 column_ids,
1335 Arc::new(AtomicU64::new(0)),
1336 ConflictBehavior::NoCheck,
1337 )
1338 .await
1339 .boxed()
1340 .execute();
1341 materialize_executor.next().await.transpose().unwrap();
1342
1343 materialize_executor.next().await.transpose().unwrap();
1344
1345 match materialize_executor.next().await.transpose().unwrap() {
1347 Some(Message::Barrier(_)) => {
1348 let row = table
1349 .get_row(
1350 &OwnedRow::new(vec![Some(3_i32.into())]),
1351 HummockReadEpoch::NoWait(u64::MAX),
1352 )
1353 .await
1354 .unwrap();
1355 assert_eq!(
1356 row,
1357 Some(OwnedRow::new(vec![Some(3_i32.into()), Some(6_i32.into())]))
1358 );
1359 }
1360 _ => unreachable!(),
1361 }
1362 materialize_executor.next().await.transpose().unwrap();
1363 match materialize_executor.next().await.transpose().unwrap() {
1365 Some(Message::Barrier(_)) => {
1366 let row = table
1367 .get_row(
1368 &OwnedRow::new(vec![Some(7_i32.into())]),
1369 HummockReadEpoch::NoWait(u64::MAX),
1370 )
1371 .await
1372 .unwrap();
1373 assert_eq!(
1374 row,
1375 Some(OwnedRow::new(vec![Some(7_i32.into()), Some(8_i32.into())]))
1376 );
1377 }
1378 _ => unreachable!(),
1379 }
1380 }
1381
1382 #[tokio::test]
1384 async fn test_upsert_stream() {
1385 let memory_state_store = MemoryStateStore::new();
1387 let table_id = TableId::new(1);
1388 let schema = Schema::new(vec![
1390 Field::unnamed(DataType::Int32),
1391 Field::unnamed(DataType::Int32),
1392 ]);
1393 let column_ids = vec![0.into(), 1.into()];
1394
1395 let chunk1 = StreamChunk::from_pretty(
1397 " i i
1398 + 1 1",
1399 );
1400
1401 let chunk2 = StreamChunk::from_pretty(
1402 " i i
1403 + 1 2
1404 - 1 2",
1405 );
1406
1407 let source = MockSource::with_messages(vec![
1409 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1410 Message::Chunk(chunk1),
1411 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1412 Message::Chunk(chunk2),
1413 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
1414 ])
1415 .into_executor(schema.clone(), StreamKey::new());
1416
1417 let order_types = vec![OrderType::ascending()];
1418 let column_descs = vec![
1419 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
1420 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
1421 ];
1422
1423 let table = BatchTable::for_test(
1424 memory_state_store.clone(),
1425 table_id,
1426 column_descs,
1427 order_types,
1428 vec![0],
1429 vec![0, 1],
1430 );
1431
1432 let mut materialize_executor = MaterializeExecutor::for_test(
1433 source,
1434 memory_state_store,
1435 table_id,
1436 vec![ColumnOrder::new(0, OrderType::ascending())],
1437 column_ids,
1438 Arc::new(AtomicU64::new(0)),
1439 ConflictBehavior::Overwrite,
1440 )
1441 .await
1442 .boxed()
1443 .execute();
1444 materialize_executor.next().await.transpose().unwrap();
1445
1446 materialize_executor.next().await.transpose().unwrap();
1447 materialize_executor.next().await.transpose().unwrap();
1448 materialize_executor.next().await.transpose().unwrap();
1449
1450 match materialize_executor.next().await.transpose().unwrap() {
1451 Some(Message::Barrier(_)) => {
1452 let row = table
1453 .get_row(
1454 &OwnedRow::new(vec![Some(1_i32.into())]),
1455 HummockReadEpoch::NoWait(u64::MAX),
1456 )
1457 .await
1458 .unwrap();
1459 assert!(row.is_none());
1460 }
1461 _ => unreachable!(),
1462 }
1463 }
1464
1465 #[tokio::test]
1466 async fn test_check_insert_conflict() {
1467 let memory_state_store = MemoryStateStore::new();
1469 let table_id = TableId::new(1);
1470 let schema = Schema::new(vec![
1472 Field::unnamed(DataType::Int32),
1473 Field::unnamed(DataType::Int32),
1474 ]);
1475 let column_ids = vec![0.into(), 1.into()];
1476
1477 let chunk1 = StreamChunk::from_pretty(
1479 " i i
1480 + 1 3
1481 + 1 4
1482 + 2 5
1483 + 3 6",
1484 );
1485
1486 let chunk2 = StreamChunk::from_pretty(
1487 " i i
1488 + 1 3
1489 + 2 6",
1490 );
1491
1492 let chunk3 = StreamChunk::from_pretty(
1494 " i i
1495 + 1 4",
1496 );
1497
1498 let source = MockSource::with_messages(vec![
1500 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1501 Message::Chunk(chunk1),
1502 Message::Chunk(chunk2),
1503 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1504 Message::Chunk(chunk3),
1505 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
1506 ])
1507 .into_executor(schema.clone(), StreamKey::new());
1508
1509 let order_types = vec![OrderType::ascending()];
1510 let column_descs = vec![
1511 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
1512 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
1513 ];
1514
1515 let table = BatchTable::for_test(
1516 memory_state_store.clone(),
1517 table_id,
1518 column_descs,
1519 order_types,
1520 vec![0],
1521 vec![0, 1],
1522 );
1523
1524 let mut materialize_executor = MaterializeExecutor::for_test(
1525 source,
1526 memory_state_store,
1527 table_id,
1528 vec![ColumnOrder::new(0, OrderType::ascending())],
1529 column_ids,
1530 Arc::new(AtomicU64::new(0)),
1531 ConflictBehavior::Overwrite,
1532 )
1533 .await
1534 .boxed()
1535 .execute();
1536 materialize_executor.next().await.transpose().unwrap();
1537
1538 materialize_executor.next().await.transpose().unwrap();
1539 materialize_executor.next().await.transpose().unwrap();
1540
1541 match materialize_executor.next().await.transpose().unwrap() {
1543 Some(Message::Barrier(_)) => {
1544 let row = table
1545 .get_row(
1546 &OwnedRow::new(vec![Some(3_i32.into())]),
1547 HummockReadEpoch::NoWait(u64::MAX),
1548 )
1549 .await
1550 .unwrap();
1551 assert_eq!(
1552 row,
1553 Some(OwnedRow::new(vec![Some(3_i32.into()), Some(6_i32.into())]))
1554 );
1555
1556 let row = table
1557 .get_row(
1558 &OwnedRow::new(vec![Some(1_i32.into())]),
1559 HummockReadEpoch::NoWait(u64::MAX),
1560 )
1561 .await
1562 .unwrap();
1563 assert_eq!(
1564 row,
1565 Some(OwnedRow::new(vec![Some(1_i32.into()), Some(3_i32.into())]))
1566 );
1567
1568 let row = table
1569 .get_row(
1570 &OwnedRow::new(vec![Some(2_i32.into())]),
1571 HummockReadEpoch::NoWait(u64::MAX),
1572 )
1573 .await
1574 .unwrap();
1575 assert_eq!(
1576 row,
1577 Some(OwnedRow::new(vec![Some(2_i32.into()), Some(6_i32.into())]))
1578 );
1579 }
1580 _ => unreachable!(),
1581 }
1582 }
1583
1584 #[tokio::test]
1585 async fn test_delete_and_update_conflict() {
1586 let memory_state_store = MemoryStateStore::new();
1588 let table_id = TableId::new(1);
1589 let schema = Schema::new(vec![
1591 Field::unnamed(DataType::Int32),
1592 Field::unnamed(DataType::Int32),
1593 ]);
1594 let column_ids = vec![0.into(), 1.into()];
1595
1596 let chunk1 = StreamChunk::from_pretty(
1598 " i i
1599 + 1 4
1600 + 2 5
1601 + 3 6
1602 U- 8 1
1603 U+ 8 2
1604 + 8 3",
1605 );
1606
1607 let chunk2 = StreamChunk::from_pretty(
1609 " i i
1610 + 7 8
1611 - 3 4
1612 - 5 0",
1613 );
1614
1615 let chunk3 = StreamChunk::from_pretty(
1617 " i i
1618 + 1 5
1619 U- 2 4
1620 U+ 2 8
1621 U- 9 0
1622 U+ 9 1",
1623 );
1624
1625 let source = MockSource::with_messages(vec![
1627 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1628 Message::Chunk(chunk1),
1629 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1630 Message::Chunk(chunk2),
1631 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
1632 Message::Chunk(chunk3),
1633 Message::Barrier(Barrier::new_test_barrier(test_epoch(4))),
1634 ])
1635 .into_executor(schema.clone(), StreamKey::new());
1636
1637 let order_types = vec![OrderType::ascending()];
1638 let column_descs = vec![
1639 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
1640 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
1641 ];
1642
1643 let table = BatchTable::for_test(
1644 memory_state_store.clone(),
1645 table_id,
1646 column_descs,
1647 order_types,
1648 vec![0],
1649 vec![0, 1],
1650 );
1651
1652 let mut materialize_executor = MaterializeExecutor::for_test(
1653 source,
1654 memory_state_store,
1655 table_id,
1656 vec![ColumnOrder::new(0, OrderType::ascending())],
1657 column_ids,
1658 Arc::new(AtomicU64::new(0)),
1659 ConflictBehavior::Overwrite,
1660 )
1661 .await
1662 .boxed()
1663 .execute();
1664 materialize_executor.next().await.transpose().unwrap();
1665
1666 materialize_executor.next().await.transpose().unwrap();
1667
1668 match materialize_executor.next().await.transpose().unwrap() {
1670 Some(Message::Barrier(_)) => {
1671 let row = table
1673 .get_row(
1674 &OwnedRow::new(vec![Some(8_i32.into())]),
1675 HummockReadEpoch::NoWait(u64::MAX),
1676 )
1677 .await
1678 .unwrap();
1679 assert_eq!(
1680 row,
1681 Some(OwnedRow::new(vec![Some(8_i32.into()), Some(3_i32.into())]))
1682 );
1683 }
1684 _ => unreachable!(),
1685 }
1686 materialize_executor.next().await.transpose().unwrap();
1687
1688 match materialize_executor.next().await.transpose().unwrap() {
1689 Some(Message::Barrier(_)) => {
1690 let row = table
1691 .get_row(
1692 &OwnedRow::new(vec![Some(7_i32.into())]),
1693 HummockReadEpoch::NoWait(u64::MAX),
1694 )
1695 .await
1696 .unwrap();
1697 assert_eq!(
1698 row,
1699 Some(OwnedRow::new(vec![Some(7_i32.into()), Some(8_i32.into())]))
1700 );
1701
1702 let row = table
1704 .get_row(
1705 &OwnedRow::new(vec![Some(3_i32.into())]),
1706 HummockReadEpoch::NoWait(u64::MAX),
1707 )
1708 .await
1709 .unwrap();
1710 assert_eq!(row, None);
1711
1712 let row = table
1714 .get_row(
1715 &OwnedRow::new(vec![Some(5_i32.into())]),
1716 HummockReadEpoch::NoWait(u64::MAX),
1717 )
1718 .await
1719 .unwrap();
1720 assert_eq!(row, None);
1721 }
1722 _ => unreachable!(),
1723 }
1724
1725 materialize_executor.next().await.transpose().unwrap();
1726 match materialize_executor.next().await.transpose().unwrap() {
1728 Some(Message::Barrier(_)) => {
1729 let row = table
1730 .get_row(
1731 &OwnedRow::new(vec![Some(1_i32.into())]),
1732 HummockReadEpoch::NoWait(u64::MAX),
1733 )
1734 .await
1735 .unwrap();
1736 assert_eq!(
1737 row,
1738 Some(OwnedRow::new(vec![Some(1_i32.into()), Some(5_i32.into())]))
1739 );
1740
1741 let row = table
1743 .get_row(
1744 &OwnedRow::new(vec![Some(2_i32.into())]),
1745 HummockReadEpoch::NoWait(u64::MAX),
1746 )
1747 .await
1748 .unwrap();
1749 assert_eq!(
1750 row,
1751 Some(OwnedRow::new(vec![Some(2_i32.into()), Some(8_i32.into())]))
1752 );
1753
1754 let row = table
1756 .get_row(
1757 &OwnedRow::new(vec![Some(9_i32.into())]),
1758 HummockReadEpoch::NoWait(u64::MAX),
1759 )
1760 .await
1761 .unwrap();
1762 assert_eq!(
1763 row,
1764 Some(OwnedRow::new(vec![Some(9_i32.into()), Some(1_i32.into())]))
1765 );
1766 }
1767 _ => unreachable!(),
1768 }
1769 }
1770
1771 #[tokio::test]
1772 async fn test_change_buffer_into_chunk_with_stream_key() {
1773 let memory_state_store = MemoryStateStore::new();
1776 let table_id = TableId::new(1);
1777 let schema = Schema::new(vec![
1778 Field::unnamed(DataType::Int32),
1779 Field::unnamed(DataType::Int32),
1780 ]);
1781 let column_ids = vec![0.into(), 1.into()];
1782
1783 let chunk1 = StreamChunk::from_pretty(
1784 " i i
1785 + 1 4",
1786 );
1787 let chunk2 = StreamChunk::from_pretty(
1788 " i i
1789 + 1 5",
1790 );
1791
1792 let source = MockSource::with_messages(vec![
1793 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1794 Message::Chunk(chunk1),
1795 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1796 Message::Chunk(chunk2),
1797 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
1798 ])
1799 .into_executor(schema.clone(), StreamKey::new());
1800
1801 let mut materialize_executor = MaterializeExecutor::for_test_with_stream_key(
1802 source,
1803 memory_state_store,
1804 table_id,
1805 vec![ColumnOrder::new(0, OrderType::ascending())],
1806 vec![0, 1],
1807 column_ids,
1808 Arc::new(AtomicU64::new(0)),
1809 ConflictBehavior::Overwrite,
1810 )
1811 .await
1812 .boxed()
1813 .execute();
1814
1815 materialize_executor.next().await.transpose().unwrap();
1817 materialize_executor.next().await.transpose().unwrap();
1818
1819 materialize_executor.next().await.transpose().unwrap();
1821
1822 match materialize_executor.next().await.transpose().unwrap() {
1824 Some(Message::Chunk(chunk)) => {
1825 assert_eq!(
1826 chunk.compact_vis(),
1827 StreamChunk::from_pretty(
1828 " i i
1829 - 1 4
1830 + 1 5"
1831 )
1832 );
1833 }
1834 other => panic!("expect chunk, got {other:?}"),
1835 }
1836 }
1837
1838 #[tokio::test]
1839 async fn test_ignore_insert_conflict() {
1840 let memory_state_store = MemoryStateStore::new();
1842 let table_id = TableId::new(1);
1843 let schema = Schema::new(vec![
1845 Field::unnamed(DataType::Int32),
1846 Field::unnamed(DataType::Int32),
1847 ]);
1848 let column_ids = vec![0.into(), 1.into()];
1849
1850 let chunk1 = StreamChunk::from_pretty(
1852 " i i
1853 + 1 3
1854 + 1 4
1855 + 2 5
1856 + 3 6",
1857 );
1858
1859 let chunk2 = StreamChunk::from_pretty(
1860 " i i
1861 + 1 5
1862 + 2 6",
1863 );
1864
1865 let chunk3 = StreamChunk::from_pretty(
1867 " i i
1868 + 1 6",
1869 );
1870
1871 let source = MockSource::with_messages(vec![
1873 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1874 Message::Chunk(chunk1),
1875 Message::Chunk(chunk2),
1876 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1877 Message::Chunk(chunk3),
1878 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
1879 ])
1880 .into_executor(schema.clone(), StreamKey::new());
1881
1882 let order_types = vec![OrderType::ascending()];
1883 let column_descs = vec![
1884 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
1885 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
1886 ];
1887
1888 let table = BatchTable::for_test(
1889 memory_state_store.clone(),
1890 table_id,
1891 column_descs,
1892 order_types,
1893 vec![0],
1894 vec![0, 1],
1895 );
1896
1897 let mut materialize_executor = MaterializeExecutor::for_test(
1898 source,
1899 memory_state_store,
1900 table_id,
1901 vec![ColumnOrder::new(0, OrderType::ascending())],
1902 column_ids,
1903 Arc::new(AtomicU64::new(0)),
1904 ConflictBehavior::IgnoreConflict,
1905 )
1906 .await
1907 .boxed()
1908 .execute();
1909 materialize_executor.next().await.transpose().unwrap();
1910
1911 materialize_executor.next().await.transpose().unwrap();
1912 materialize_executor.next().await.transpose().unwrap();
1913
1914 match materialize_executor.next().await.transpose().unwrap() {
1916 Some(Message::Barrier(_)) => {
1917 let row = table
1918 .get_row(
1919 &OwnedRow::new(vec![Some(3_i32.into())]),
1920 HummockReadEpoch::NoWait(u64::MAX),
1921 )
1922 .await
1923 .unwrap();
1924 assert_eq!(
1925 row,
1926 Some(OwnedRow::new(vec![Some(3_i32.into()), Some(6_i32.into())]))
1927 );
1928
1929 let row = table
1930 .get_row(
1931 &OwnedRow::new(vec![Some(1_i32.into())]),
1932 HummockReadEpoch::NoWait(u64::MAX),
1933 )
1934 .await
1935 .unwrap();
1936 assert_eq!(
1937 row,
1938 Some(OwnedRow::new(vec![Some(1_i32.into()), Some(3_i32.into())]))
1939 );
1940
1941 let row = table
1942 .get_row(
1943 &OwnedRow::new(vec![Some(2_i32.into())]),
1944 HummockReadEpoch::NoWait(u64::MAX),
1945 )
1946 .await
1947 .unwrap();
1948 assert_eq!(
1949 row,
1950 Some(OwnedRow::new(vec![Some(2_i32.into()), Some(5_i32.into())]))
1951 );
1952 }
1953 _ => unreachable!(),
1954 }
1955 }
1956
1957 #[tokio::test]
1958 async fn test_ignore_delete_then_insert() {
1959 let memory_state_store = MemoryStateStore::new();
1961 let table_id = TableId::new(1);
1962 let schema = Schema::new(vec![
1964 Field::unnamed(DataType::Int32),
1965 Field::unnamed(DataType::Int32),
1966 ]);
1967 let column_ids = vec![0.into(), 1.into()];
1968
1969 let chunk1 = StreamChunk::from_pretty(
1971 " i i
1972 + 1 3
1973 - 1 3
1974 + 1 6",
1975 );
1976
1977 let source = MockSource::with_messages(vec![
1979 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
1980 Message::Chunk(chunk1),
1981 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
1982 ])
1983 .into_executor(schema.clone(), StreamKey::new());
1984
1985 let order_types = vec![OrderType::ascending()];
1986 let column_descs = vec![
1987 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
1988 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
1989 ];
1990
1991 let table = BatchTable::for_test(
1992 memory_state_store.clone(),
1993 table_id,
1994 column_descs,
1995 order_types,
1996 vec![0],
1997 vec![0, 1],
1998 );
1999
2000 let mut materialize_executor = MaterializeExecutor::for_test(
2001 source,
2002 memory_state_store,
2003 table_id,
2004 vec![ColumnOrder::new(0, OrderType::ascending())],
2005 column_ids,
2006 Arc::new(AtomicU64::new(0)),
2007 ConflictBehavior::IgnoreConflict,
2008 )
2009 .await
2010 .boxed()
2011 .execute();
2012 let _msg1 = materialize_executor
2013 .next()
2014 .await
2015 .transpose()
2016 .unwrap()
2017 .unwrap()
2018 .as_barrier()
2019 .unwrap();
2020 let _msg2 = materialize_executor
2021 .next()
2022 .await
2023 .transpose()
2024 .unwrap()
2025 .unwrap()
2026 .as_chunk()
2027 .unwrap();
2028 let _msg3 = materialize_executor
2029 .next()
2030 .await
2031 .transpose()
2032 .unwrap()
2033 .unwrap()
2034 .as_barrier()
2035 .unwrap();
2036
2037 let row = table
2038 .get_row(
2039 &OwnedRow::new(vec![Some(1_i32.into())]),
2040 HummockReadEpoch::NoWait(u64::MAX),
2041 )
2042 .await
2043 .unwrap();
2044 assert_eq!(
2045 row,
2046 Some(OwnedRow::new(vec![Some(1_i32.into()), Some(6_i32.into())]))
2047 );
2048 }
2049
2050 #[tokio::test]
2051 async fn test_ignore_delete_and_update_conflict() {
2052 let memory_state_store = MemoryStateStore::new();
2054 let table_id = TableId::new(1);
2055 let schema = Schema::new(vec![
2057 Field::unnamed(DataType::Int32),
2058 Field::unnamed(DataType::Int32),
2059 ]);
2060 let column_ids = vec![0.into(), 1.into()];
2061
2062 let chunk1 = StreamChunk::from_pretty(
2064 " i i
2065 + 1 4
2066 + 2 5
2067 + 3 6
2068 U- 8 1
2069 U+ 8 2
2070 + 8 3",
2071 );
2072
2073 let chunk2 = StreamChunk::from_pretty(
2075 " i i
2076 + 7 8
2077 - 3 4
2078 - 5 0",
2079 );
2080
2081 let chunk3 = StreamChunk::from_pretty(
2083 " i i
2084 + 1 5
2085 U- 2 4
2086 U+ 2 8
2087 U- 9 0
2088 U+ 9 1",
2089 );
2090
2091 let source = MockSource::with_messages(vec![
2093 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
2094 Message::Chunk(chunk1),
2095 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
2096 Message::Chunk(chunk2),
2097 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
2098 Message::Chunk(chunk3),
2099 Message::Barrier(Barrier::new_test_barrier(test_epoch(4))),
2100 ])
2101 .into_executor(schema.clone(), StreamKey::new());
2102
2103 let order_types = vec![OrderType::ascending()];
2104 let column_descs = vec![
2105 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
2106 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
2107 ];
2108
2109 let table = BatchTable::for_test(
2110 memory_state_store.clone(),
2111 table_id,
2112 column_descs,
2113 order_types,
2114 vec![0],
2115 vec![0, 1],
2116 );
2117
2118 let mut materialize_executor = MaterializeExecutor::for_test(
2119 source,
2120 memory_state_store,
2121 table_id,
2122 vec![ColumnOrder::new(0, OrderType::ascending())],
2123 column_ids,
2124 Arc::new(AtomicU64::new(0)),
2125 ConflictBehavior::IgnoreConflict,
2126 )
2127 .await
2128 .boxed()
2129 .execute();
2130 materialize_executor.next().await.transpose().unwrap();
2131
2132 materialize_executor.next().await.transpose().unwrap();
2133
2134 match materialize_executor.next().await.transpose().unwrap() {
2136 Some(Message::Barrier(_)) => {
2137 let row = table
2139 .get_row(
2140 &OwnedRow::new(vec![Some(8_i32.into())]),
2141 HummockReadEpoch::NoWait(u64::MAX),
2142 )
2143 .await
2144 .unwrap();
2145 assert_eq!(
2146 row,
2147 Some(OwnedRow::new(vec![Some(8_i32.into()), Some(2_i32.into())]))
2148 );
2149 }
2150 _ => unreachable!(),
2151 }
2152 materialize_executor.next().await.transpose().unwrap();
2153
2154 match materialize_executor.next().await.transpose().unwrap() {
2155 Some(Message::Barrier(_)) => {
2156 let row = table
2157 .get_row(
2158 &OwnedRow::new(vec![Some(7_i32.into())]),
2159 HummockReadEpoch::NoWait(u64::MAX),
2160 )
2161 .await
2162 .unwrap();
2163 assert_eq!(
2164 row,
2165 Some(OwnedRow::new(vec![Some(7_i32.into()), Some(8_i32.into())]))
2166 );
2167
2168 let row = table
2170 .get_row(
2171 &OwnedRow::new(vec![Some(3_i32.into())]),
2172 HummockReadEpoch::NoWait(u64::MAX),
2173 )
2174 .await
2175 .unwrap();
2176 assert_eq!(row, None);
2177
2178 let row = table
2180 .get_row(
2181 &OwnedRow::new(vec![Some(5_i32.into())]),
2182 HummockReadEpoch::NoWait(u64::MAX),
2183 )
2184 .await
2185 .unwrap();
2186 assert_eq!(row, None);
2187 }
2188 _ => unreachable!(),
2189 }
2190
2191 materialize_executor.next().await.transpose().unwrap();
2192 match materialize_executor.next().await.transpose().unwrap() {
2195 Some(Message::Barrier(_)) => {
2196 let row = table
2197 .get_row(
2198 &OwnedRow::new(vec![Some(1_i32.into())]),
2199 HummockReadEpoch::NoWait(u64::MAX),
2200 )
2201 .await
2202 .unwrap();
2203 assert_eq!(
2204 row,
2205 Some(OwnedRow::new(vec![Some(1_i32.into()), Some(4_i32.into())]))
2206 );
2207
2208 let row = table
2210 .get_row(
2211 &OwnedRow::new(vec![Some(2_i32.into())]),
2212 HummockReadEpoch::NoWait(u64::MAX),
2213 )
2214 .await
2215 .unwrap();
2216 assert_eq!(
2217 row,
2218 Some(OwnedRow::new(vec![Some(2_i32.into()), Some(8_i32.into())]))
2219 );
2220
2221 let row = table
2223 .get_row(
2224 &OwnedRow::new(vec![Some(9_i32.into())]),
2225 HummockReadEpoch::NoWait(u64::MAX),
2226 )
2227 .await
2228 .unwrap();
2229 assert_eq!(
2230 row,
2231 Some(OwnedRow::new(vec![Some(9_i32.into()), Some(1_i32.into())]))
2232 );
2233 }
2234 _ => unreachable!(),
2235 }
2236 }
2237
2238 #[tokio::test]
2239 async fn test_do_update_if_not_null_conflict() {
2240 let memory_state_store = MemoryStateStore::new();
2242 let table_id = TableId::new(1);
2243 let schema = Schema::new(vec![
2245 Field::unnamed(DataType::Int32),
2246 Field::unnamed(DataType::Int32),
2247 ]);
2248 let column_ids = vec![0.into(), 1.into()];
2249
2250 let chunk1 = StreamChunk::from_pretty(
2252 " i i
2253 + 1 4
2254 + 2 .
2255 + 3 6
2256 U- 8 .
2257 U+ 8 2
2258 + 8 .",
2259 );
2260
2261 let chunk2 = StreamChunk::from_pretty(
2263 " i i
2264 + 7 8
2265 - 3 4
2266 - 5 0",
2267 );
2268
2269 let chunk3 = StreamChunk::from_pretty(
2271 " i i
2272 + 1 5
2273 + 7 .
2274 U- 2 4
2275 U+ 2 .
2276 U- 9 0
2277 U+ 9 1",
2278 );
2279
2280 let source = MockSource::with_messages(vec![
2282 Message::Barrier(Barrier::new_test_barrier(test_epoch(1))),
2283 Message::Chunk(chunk1),
2284 Message::Barrier(Barrier::new_test_barrier(test_epoch(2))),
2285 Message::Chunk(chunk2),
2286 Message::Barrier(Barrier::new_test_barrier(test_epoch(3))),
2287 Message::Chunk(chunk3),
2288 Message::Barrier(Barrier::new_test_barrier(test_epoch(4))),
2289 ])
2290 .into_executor(schema.clone(), StreamKey::new());
2291
2292 let order_types = vec![OrderType::ascending()];
2293 let column_descs = vec![
2294 ColumnDesc::unnamed(column_ids[0], DataType::Int32),
2295 ColumnDesc::unnamed(column_ids[1], DataType::Int32),
2296 ];
2297
2298 let table = BatchTable::for_test(
2299 memory_state_store.clone(),
2300 table_id,
2301 column_descs,
2302 order_types,
2303 vec![0],
2304 vec![0, 1],
2305 );
2306
2307 let mut materialize_executor = MaterializeExecutor::for_test(
2308 source,
2309 memory_state_store,
2310 table_id,
2311 vec![ColumnOrder::new(0, OrderType::ascending())],
2312 column_ids,
2313 Arc::new(AtomicU64::new(0)),
2314 ConflictBehavior::DoUpdateIfNotNull,
2315 )
2316 .await
2317 .boxed()
2318 .execute();
2319 materialize_executor.next().await.transpose().unwrap();
2320
2321 materialize_executor.next().await.transpose().unwrap();
2322
2323 match materialize_executor.next().await.transpose().unwrap() {
2325 Some(Message::Barrier(_)) => {
2326 let row = table
2327 .get_row(
2328 &OwnedRow::new(vec![Some(8_i32.into())]),
2329 HummockReadEpoch::NoWait(u64::MAX),
2330 )
2331 .await
2332 .unwrap();
2333 assert_eq!(
2334 row,
2335 Some(OwnedRow::new(vec![Some(8_i32.into()), Some(2_i32.into())]))
2336 );
2337
2338 let row = table
2339 .get_row(
2340 &OwnedRow::new(vec![Some(2_i32.into())]),
2341 HummockReadEpoch::NoWait(u64::MAX),
2342 )
2343 .await
2344 .unwrap();
2345 assert_eq!(row, Some(OwnedRow::new(vec![Some(2_i32.into()), None])));
2346 }
2347 _ => unreachable!(),
2348 }
2349 materialize_executor.next().await.transpose().unwrap();
2350
2351 match materialize_executor.next().await.transpose().unwrap() {
2352 Some(Message::Barrier(_)) => {
2353 let row = table
2354 .get_row(
2355 &OwnedRow::new(vec![Some(7_i32.into())]),
2356 HummockReadEpoch::NoWait(u64::MAX),
2357 )
2358 .await
2359 .unwrap();
2360 assert_eq!(
2361 row,
2362 Some(OwnedRow::new(vec![Some(7_i32.into()), Some(8_i32.into())]))
2363 );
2364
2365 let row = table
2367 .get_row(
2368 &OwnedRow::new(vec![Some(3_i32.into())]),
2369 HummockReadEpoch::NoWait(u64::MAX),
2370 )
2371 .await
2372 .unwrap();
2373 assert_eq!(row, None);
2374
2375 let row = table
2377 .get_row(
2378 &OwnedRow::new(vec![Some(5_i32.into())]),
2379 HummockReadEpoch::NoWait(u64::MAX),
2380 )
2381 .await
2382 .unwrap();
2383 assert_eq!(row, None);
2384 }
2385 _ => unreachable!(),
2386 }
2387
2388 materialize_executor.next().await.transpose().unwrap();
2389 match materialize_executor.next().await.transpose().unwrap() {
2392 Some(Message::Barrier(_)) => {
2393 let row = table
2394 .get_row(
2395 &OwnedRow::new(vec![Some(7_i32.into())]),
2396 HummockReadEpoch::NoWait(u64::MAX),
2397 )
2398 .await
2399 .unwrap();
2400 assert_eq!(
2401 row,
2402 Some(OwnedRow::new(vec![Some(7_i32.into()), Some(8_i32.into())]))
2403 );
2404
2405 let row = table
2407 .get_row(
2408 &OwnedRow::new(vec![Some(2_i32.into())]),
2409 HummockReadEpoch::NoWait(u64::MAX),
2410 )
2411 .await
2412 .unwrap();
2413 assert_eq!(row, Some(OwnedRow::new(vec![Some(2_i32.into()), None])));
2414
2415 let row = table
2417 .get_row(
2418 &OwnedRow::new(vec![Some(9_i32.into())]),
2419 HummockReadEpoch::NoWait(u64::MAX),
2420 )
2421 .await
2422 .unwrap();
2423 assert_eq!(
2424 row,
2425 Some(OwnedRow::new(vec![Some(9_i32.into()), Some(1_i32.into())]))
2426 );
2427 }
2428 _ => unreachable!(),
2429 }
2430 }
2431
2432 fn gen_fuzz_data(row_number: usize, chunk_size: usize) -> Vec<StreamChunk> {
2433 const KN: u32 = 4;
2434 const SEED: u64 = 998244353;
2435 let mut ret = vec![];
2436 let mut builder =
2437 StreamChunkBuilder::new(chunk_size, vec![DataType::Int32, DataType::Int32]);
2438 let mut rng = SmallRng::seed_from_u64(SEED);
2439
2440 let random_vis = |c: StreamChunk, rng: &mut SmallRng| -> StreamChunk {
2441 let len = c.data_chunk().capacity();
2442 let mut c = StreamChunkMut::from(c);
2443 for i in 0..len {
2444 c.set_vis(i, rng.random_bool(0.5));
2445 }
2446 c.into()
2447 };
2448 for _ in 0..row_number {
2449 let k = (rng.next_u32() % KN) as i32;
2450 let v = rng.next_u32() as i32;
2451 let op = if rng.random_bool(0.5) {
2452 Op::Insert
2453 } else {
2454 Op::Delete
2455 };
2456 if let Some(c) =
2457 builder.append_row(op, OwnedRow::new(vec![Some(k.into()), Some(v.into())]))
2458 {
2459 ret.push(random_vis(c, &mut rng));
2460 }
2461 }
2462 if let Some(c) = builder.take() {
2463 ret.push(random_vis(c, &mut rng));
2464 }
2465 ret
2466 }
2467
2468 async fn fuzz_test_stream_consistent_inner(conflict_behavior: ConflictBehavior) {
2469 const N: usize = 100000;
2470
2471 let memory_state_store = MemoryStateStore::new();
2473 let table_id = TableId::new(1);
2474 let schema = Schema::new(vec![
2476 Field::unnamed(DataType::Int32),
2477 Field::unnamed(DataType::Int32),
2478 ]);
2479 let column_ids = vec![0.into(), 1.into()];
2480
2481 let chunks = gen_fuzz_data(N, 128);
2482 let messages = iter::once(Message::Barrier(Barrier::new_test_barrier(test_epoch(1))))
2483 .chain(chunks.into_iter().map(Message::Chunk))
2484 .chain(iter::once(Message::Barrier(Barrier::new_test_barrier(
2485 test_epoch(2),
2486 ))))
2487 .collect();
2488 let source =
2490 MockSource::with_messages(messages).into_executor(schema.clone(), StreamKey::new());
2491
2492 let mut materialize_executor = MaterializeExecutor::for_test(
2493 source,
2494 memory_state_store.clone(),
2495 table_id,
2496 vec![ColumnOrder::new(0, OrderType::ascending())],
2497 column_ids,
2498 Arc::new(AtomicU64::new(0)),
2499 conflict_behavior,
2500 )
2501 .await
2502 .boxed()
2503 .execute();
2504 materialize_executor.expect_barrier().await;
2505
2506 let order_types = vec![OrderType::ascending()];
2507 let column_descs = vec![
2508 ColumnDesc::unnamed(0.into(), DataType::Int32),
2509 ColumnDesc::unnamed(1.into(), DataType::Int32),
2510 ];
2511 let pk_indices = vec![0];
2512
2513 let mut table = StateTable::from_table_catalog(
2514 &crate::common::table::test_utils::gen_pbtable(
2515 TableId::from(1002),
2516 column_descs.clone(),
2517 order_types,
2518 pk_indices,
2519 0,
2520 ),
2521 memory_state_store.clone(),
2522 None,
2523 )
2524 .await;
2525
2526 while let Message::Chunk(c) = materialize_executor.next().await.unwrap().unwrap() {
2527 table.write_chunk(c);
2529 }
2530 }
2531
2532 #[tokio::test]
2533 async fn fuzz_test_stream_consistent_upsert() {
2534 fuzz_test_stream_consistent_inner(ConflictBehavior::Overwrite).await
2535 }
2536
2537 #[tokio::test]
2538 async fn fuzz_test_stream_consistent_ignore() {
2539 fuzz_test_stream_consistent_inner(ConflictBehavior::IgnoreConflict).await
2540 }
2541}