Skip to main content

risingwave_stream/executor/mview/
materialize.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::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
69/// `MaterializeExecutor` materializes changes in stream into a materialized view on storage.
70pub 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 of the *output* of this materialize node.
78    ///
79    /// This can be different from the table PK. Typically, there are 3 cases:
80    ///
81    /// - Normal `TABLE`: `stream_key == pk`, so no special handling is needed.
82    /// - TTL-ed `TABLE` (`WITH TTL`): `stream_key` is a superset of `pk` by appending the TTL
83    ///   watermark column. In this case, updates to the same pk may have different stream keys, so
84    ///   we must rewrite such `Update` into `Delete + Insert` before yielding.
85    /// - `MV` or `INDEX`: `pk` can be a superset of `stream_key` to also include the order key or
86    ///   distribution key specified by the user. No special handling is needed either.
87    stream_key_indices: Vec<usize>,
88
89    /// Columns of arrange keys (including pk, group keys, join keys, etc.)
90    arrange_key_indices: Vec<usize>,
91
92    actor_context: ActorContextRef,
93
94    /// The cache for conflict handling. `None` if conflict behavior is `NoCheck`.
95    materialize_cache: Option<MaterializeCache>,
96
97    /// Indices of columns that may carry a Debezium unchanged-TOAST placeholder
98    /// (PostgreSQL CDC only). Mirrors the value passed into `MaterializeCache`.
99    /// Used at runtime to prevent the `Overwrite -> NoCheck` downgrade from
100    /// bypassing the placeholder replacement path.
101    toastable_column_indices: Option<Vec<usize>>,
102
103    conflict_behavior: ConflictBehavior,
104
105    /// Whether the table can clean itself by TTL watermark, i.e., is defined with `WATERMARK ... WITH TTL`.
106    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    /// No data will be written to hummock table. This Materialize is just a dummy node.
117    /// Used for APPEND ONLY table with iceberg engine. All data will be written to iceberg table directly.
118    is_dummy_table: bool,
119
120    /// Optional refresh arguments and state for refreshable materialized views
121    refresh_args: Option<RefreshableMaterializeArgs<S, SD>>,
122
123    /// Local barrier manager for reporting barrier events
124    local_barrier_manager: LocalBarrierManager,
125}
126
127/// Arguments and state for refreshable materialized views
128pub struct RefreshableMaterializeArgs<S: StateStore, SD: ValueRowSerde> {
129    /// Table catalog for main table
130    pub table_catalog: Table,
131
132    /// Table catalog for staging table
133    pub staging_table_catalog: Table,
134
135    /// Flag indicating if this table is currently being refreshed
136    pub is_refreshing: bool,
137
138    /// During data refresh (between `RefreshStart` and `LoadFinish`),
139    /// data will be written to both the main table and the staging table.
140    ///
141    /// The staging table is PK-only.
142    ///
143    /// After `LoadFinish`, we will do a `DELETE FROM main_table WHERE pk NOT IN (SELECT pk FROM staging_table)`, and then purge the staging table.
144    pub staging_table: StateTableInner<S, SD>,
145
146    /// Progress table for tracking refresh state per `VNode` for fault tolerance
147    pub progress_table: RefreshProgressTable<S>,
148
149    /// Table ID for this refreshable materialized view
150    pub table_id: TableId,
151}
152
153impl<S: StateStore, SD: ValueRowSerde> RefreshableMaterializeArgs<S, SD> {
154    /// Create new `RefreshableMaterializeArgs`
155    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        // staging table is pk-only, and we don't need to check value consistency
165        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        // Get primary key length from main table catalog
180        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        // For tables with watermark TTL. Due to async state cleaning, it's uncertain whether an expired
204        // key has been cleaned up by the storage or not, thus we are not sure if to write `Update` or `Insert`
205        // if the key appears again.
206        StateTableOpConsistencyLevel::Inconsistent
207    } else if !subscriber_ids.is_empty() {
208        StateTableOpConsistencyLevel::LogStoreEnabled
209    } else if !may_have_downstream && matches!(conflict_behavior, ConflictBehavior::Overwrite) {
210        // Table with overwrite conflict behavior could disable conflict check
211        // if no downstream mv depends on it, so we use a inconsistent_op to skip sanity check as well.
212        StateTableOpConsistencyLevel::Inconsistent
213    } else {
214        StateTableOpConsistencyLevel::ConsistentOldValue
215    }
216}
217
218impl<S: StateStore, SD: ValueRowSerde> MaterializeExecutor<S, SD> {
219    /// Create a new `MaterializeExecutor` with distribution specified with `distribution_keys` and
220    /// `vnodes`. For singleton distribution, `distribution_keys` should be empty and `vnodes`
221    /// should be `None`.
222    #[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        // Extract TOAST-able column indices from table columns.
246        // Only for PostgreSQL CDC tables.
247        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                    // Currently supports TOAST updates for:
255                    // - jsonb (DataType::Jsonb)
256                    // - varchar (DataType::Varchar)
257                    // - bytea (DataType::Bytea)
258                    // - pgvector (DataType::Vector)
259                    // - One-dimensional arrays of the above types (DataType::List)
260                    //   Note: Some array types may not be fully supported yet, see issue  https://github.com/risingwavelabs/risingwave/issues/22916 for details.
261
262                    // For details on how TOAST values are handled, see comments in `is_debezium_unavailable_value`.
263                    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        // Note: The current implementation could potentially trigger a switch on the inconsistent_op flag. If the storage relies on this flag to perform optimizations, it would be advisable to maintain consistency with it throughout the lifecycle.
306        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; // Save epoch for later use (unused in normal execution)
375        // The first barrier message should be propagated.
376        yield Message::Barrier(barrier);
377        self.state_table.init_epoch(first_epoch).await?;
378
379        // default to normal ingestion
380        let mut inner_state =
381            Box::new(MaterializeStreamState::<BarrierMutationType>::NormalIngestion);
382        // Initialize staging table for refreshable materialized views
383        if let Some(ref mut refresh_args) = self.refresh_args {
384            refresh_args.staging_table.init_epoch(first_epoch).await?;
385
386            // Initialize progress table and load existing progress for recovery
387            refresh_args.progress_table.recover(first_epoch).await?;
388
389            // Check if refresh is already in progress (recovery scenario)
390            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                // Since stage info is no longer stored in progress table,
400                // we need to determine recovery state differently.
401                // For now, assume all incomplete VNodes need to continue merging
402                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                    // Some VNodes are incomplete, need to resume refresh operation
412                    tracing::info!(
413                        incomplete_vnodes = incomplete_vnodes.len(),
414                        "Recovery detected incomplete VNodes, resuming refresh operation"
415                    );
416                    // Since stage tracking is now in memory, we'll determine the appropriate
417                    // stage based on the executor's internal state machine
418                } else {
419                    // This should not happen if is_complete() returned false, but handle it gracefully
420                    tracing::warn!("Unexpected recovery state: no incomplete VNodes found");
421                }
422            }
423        }
424
425        // Determine initial execution stage (for recovery scenarios)
426        if let Some(ref refresh_args) = self.refresh_args
427            && refresh_args.is_refreshing
428        {
429            // Recovery logic: Check if there are incomplete vnodes from previous run
430            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                // Resume from merge stage since some VNodes were left incomplete
439                *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 execution loop: cycles through Stage 1 -> Stage 2 -> Stage 3 -> Stage 1...
448        '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                                // This is an optimization that handles conflicts only when a particular materialized view downstream has no MV dependencies.
479                                // This optimization is applied only when there is no specified version column and the is_consistent_op flag of the state table is false,
480                                // and the conflict behavior is overwrite. We can rely on the state table to overwrite the conflicting rows in the storage,
481                                // while outputting inconsistent changes to downstream which no one will subscribe to.
482                                // For tables with watermark TTL (indicated by `cleaned_by_ttl_watermark`), conflict check must be enabled.
483                                // TODO(ttl): differentiate between table consistency and downstream consistency.
484                                // When the table carries TOAST-able columns (PostgreSQL CDC), the
485                                // `MaterializeCache` is the only place where the Debezium
486                                // unchanged-TOAST placeholder is detected and swapped for the row's
487                                // previous value. Demoting to `NoCheck` here would skip the cache
488                                // entirely and let the placeholder land in state. Keep the original
489                                // `Overwrite` behavior in that case so the replacement path runs.
490                                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                                            // empty chunk
506                                            continue;
507                                        }
508
509                                        // For refreshable materialized views, write to staging table during refresh
510                                        // Do not use generate_output here.
511                                        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                                                // refreshable source should be append-only
524                                                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                                            // We only hit this branch for TTL-ed tables for now.
549                                            // Assert stream key is a superset of pk.
550                                            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                                        // For refreshable materialized views, also write to staging table during refresh
576                                        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                                                // refreshable source should be append-only
589                                                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                    // Scope to limit immutable borrows to state tables
646                    {
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                        // Prefer to select input stream to handle barriers promptly
658                        // Rebuild the merge stream each time processing a barrier
659                        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 stream finished
687                                            merge_complete = true;
688
689                                            // If the merge stream finished, we need to wait for the next barrier to commit states
690                                        }
691                                    }
692                                }
693                            }
694                        }
695                    }
696
697                    // Process collected rows for deletion
698                    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                    // should wait for at least one barrier
717                    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                    // Wait for staging table truncation to complete
792                    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                                // Initialize progress tracking for all VNodes
828                                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                                // Get associated source id from table catalog
838                                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                    // ===== normal operation =====
859
860                    // If a downstream mv depends on the current table, we need to do conflict check again.
861                    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                    // Commit staging table for refreshable materialized views
885                    let refresh_post_commit = if let Some(ref mut refresh_args) = self.refresh_args
886                    {
887                        // Commit progress table for fault tolerance
888
889                        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                    // Update the vnode bitmap for the state table if asked.
901                    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                    // Handle staging table post commit
911                    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                    // ====== transition to next state ======
925
926                    *inner_state = *expect_next_state;
927                }
928            }
929        }
930    }
931
932    /// Stream that yields rows to be deleted from main table.
933    /// Yields `Some((vnode, row))` for rows that exist in main but not in staging.
934    /// Yields `None` when finished processing all vnodes.
935    #[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            // Check if this VNode has already been completed (for fault tolerance)
944            let pk_range: (Bound<OwnedRow>, Bound<OwnedRow>) =
945                if let Some(current_entry) = progress_table.get_progress(vnode) {
946                    // Skip already completed VNodes during recovery
947                    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) = &current_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            // Sort-merge join implementation using dual pointers
985            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                // Advance staging iterator until we find a key >= main_key
993                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                            // main_key > staging_key, advance staging
999                            staging_item = iter_staging.next().await.transpose()?;
1000                        }
1001                        std::cmp::Ordering::Equal => {
1002                            // Keys match, this row exists in both tables, no need to delete
1003                            break;
1004                        }
1005                        std::cmp::Ordering::Less => {
1006                            // main_key < staging_key, main row doesn't exist in staging, delete it
1007                            should_delete = true;
1008                            break;
1009                        }
1010                    }
1011                }
1012
1013                // If staging_item is None, all remaining main rows should be deleted
1014                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                // Advance main iterator
1023                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            // Mark this VNode as completed
1044            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, // completed
1049                    current_entry.processed_rows,
1050                )?;
1051
1052                tracing::debug!(vnode = vnode.to_index(), "Completed merging VNode");
1053            }
1054        }
1055
1056        // Signal completion
1057        yield None;
1058    }
1059
1060    /// return true when changed
1061    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    /// Initialize refresh progress tracking for all `VNodes`
1098    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        // Initialize progress for all VNodes in the current bitmap
1106        for vnode in state_table.vnodes().iter_vnodes() {
1107            progress_table.set_progress(
1108                vnode, None,  // initial position
1109                false, // not completed yet
1110                0,     // initial processed rows
1111            )?;
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    /// Create a new `MaterializeExecutor` without distribution info for test purpose.
1125    #[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, // Test constructor doesn't support refresh functionality
1240            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        // Prepare storage and memtable.
1282        let memory_state_store = MemoryStateStore::new();
1283        let table_id = TableId::new(1);
1284        // Two columns of int32 type, the first column is PK.
1285        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        // Prepare source chunks.
1292        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        // Prepare stream executors.
1305        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        // First stream chunk. We check the existence of (3) -> (3,6)
1346        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        // Second stream chunk. We check the existence of (7) -> (7,8)
1364        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    // https://github.com/risingwavelabs/risingwave/issues/13346
1383    #[tokio::test]
1384    async fn test_upsert_stream() {
1385        // Prepare storage and memtable.
1386        let memory_state_store = MemoryStateStore::new();
1387        let table_id = TableId::new(1);
1388        // Two columns of int32 type, the first column is PK.
1389        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        // test double insert one pk, the latter needs to override the former.
1396        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        // Prepare stream executors.
1408        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        // Prepare storage and memtable.
1468        let memory_state_store = MemoryStateStore::new();
1469        let table_id = TableId::new(1);
1470        // Two columns of int32 type, the first column is PK.
1471        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        // test double insert one pk, the latter needs to override the former.
1478        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        // test delete wrong value, delete inexistent pk
1493        let chunk3 = StreamChunk::from_pretty(
1494            " i i
1495            + 1 4",
1496        );
1497
1498        // Prepare stream executors.
1499        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        // First stream chunk. We check the existence of (3) -> (3,6)
1542        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        // Prepare storage and memtable.
1587        let memory_state_store = MemoryStateStore::new();
1588        let table_id = TableId::new(1);
1589        // Two columns of int32 type, the first column is PK.
1590        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        // test double insert one pk, the latter needs to override the former.
1597        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        // test delete wrong value, delete inexistent pk
1608        let chunk2 = StreamChunk::from_pretty(
1609            " i i
1610            + 7 8
1611            - 3 4
1612            - 5 0",
1613        );
1614
1615        // test delete wrong value, delete inexistent pk
1616        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        // Prepare stream executors.
1626        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        // First stream chunk. We check the existence of (3) -> (3,6)
1669        match materialize_executor.next().await.transpose().unwrap() {
1670            Some(Message::Barrier(_)) => {
1671                // can read (8, 3), check insert after update
1672                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                // check delete wrong value
1703                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                // check delete wrong pk
1713                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        // Second stream chunk. We check the existence of (7) -> (7,8)
1727        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                // check update wrong value
1742                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                // check update wrong pk, should become insert
1755                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        // Table PK is (col0), but stream key is (col0, col1). When the value column changes due
1774        // to overwrite conflict handling, we should output `- old` + `+ new` rather than `U-`/`U+`.
1775        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        // init barrier + first insert
1816        materialize_executor.next().await.transpose().unwrap();
1817        materialize_executor.next().await.transpose().unwrap();
1818
1819        // commit barrier
1820        materialize_executor.next().await.transpose().unwrap();
1821
1822        // overwrite conflict should be converted into delete+insert due to stream key mismatch
1823        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        // Prepare storage and memtable.
1841        let memory_state_store = MemoryStateStore::new();
1842        let table_id = TableId::new(1);
1843        // Two columns of int32 type, the first column is PK.
1844        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        // test double insert one pk, the latter needs to be ignored.
1851        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        // test delete wrong value, delete inexistent pk
1866        let chunk3 = StreamChunk::from_pretty(
1867            " i i
1868            + 1 6",
1869        );
1870
1871        // Prepare stream executors.
1872        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        // First stream chunk. We check the existence of (3) -> (3,6)
1915        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        // Prepare storage and memtable.
1960        let memory_state_store = MemoryStateStore::new();
1961        let table_id = TableId::new(1);
1962        // Two columns of int32 type, the first column is PK.
1963        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        // test insert after delete one pk, the latter insert should succeed.
1970        let chunk1 = StreamChunk::from_pretty(
1971            " i i
1972            + 1 3
1973            - 1 3
1974            + 1 6",
1975        );
1976
1977        // Prepare stream executors.
1978        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        // Prepare storage and memtable.
2053        let memory_state_store = MemoryStateStore::new();
2054        let table_id = TableId::new(1);
2055        // Two columns of int32 type, the first column is PK.
2056        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        // test double insert one pk, the latter should be ignored.
2063        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        // test delete wrong value, delete inexistent pk
2074        let chunk2 = StreamChunk::from_pretty(
2075            " i i
2076            + 7 8
2077            - 3 4
2078            - 5 0",
2079        );
2080
2081        // test delete wrong value, delete inexistent pk
2082        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        // Prepare stream executors.
2092        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        // First stream chunk. We check the existence of (3) -> (3,6)
2135        match materialize_executor.next().await.transpose().unwrap() {
2136            Some(Message::Barrier(_)) => {
2137                // can read (8, 2), check insert after update
2138                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                // check delete wrong value
2169                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                // check delete wrong pk
2179                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        // materialize_executor.next().await.transpose().unwrap();
2193        // Second stream chunk. We check the existence of (7) -> (7,8)
2194        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                // check update wrong value
2209                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                // check update wrong pk, should become insert
2222                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        // Prepare storage and memtable.
2241        let memory_state_store = MemoryStateStore::new();
2242        let table_id = TableId::new(1);
2243        // Two columns of int32 type, the first column is PK.
2244        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        // should get (8, 2)
2251        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        // should not get (3, x), should not get (5, 0)
2262        let chunk2 = StreamChunk::from_pretty(
2263            " i i
2264            + 7 8
2265            - 3 4
2266            - 5 0",
2267        );
2268
2269        // should get (2, None), (7, 8)
2270        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        // Prepare stream executors.
2281        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        // First stream chunk. We check the existence of (3) -> (3,6)
2324        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                // check delete wrong value
2366                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                // check delete wrong pk
2376                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        // materialize_executor.next().await.transpose().unwrap();
2390        // Second stream chunk. We check the existence of (7) -> (7,8)
2391        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                // check update wrong value
2406                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                // check update wrong pk, should become insert
2416                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        // Prepare storage and memtable.
2472        let memory_state_store = MemoryStateStore::new();
2473        let table_id = TableId::new(1);
2474        // Two columns of int32 type, the first column is PK.
2475        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        // Prepare stream executors.
2489        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            // check with state table's memtable
2528            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}