Skip to main content

risingwave_stream/executor/backfill/cdc/
cdc_backfill.rs

1// Copyright 2023 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::BTreeMap;
16use std::future::Future;
17use std::pin::Pin;
18
19use either::Either;
20use futures::stream::select_with_strategy;
21use futures::{Stream, stream};
22use itertools::Itertools;
23use risingwave_common::array::{DataChunk, Op};
24use risingwave_common::bail;
25use risingwave_common::bitmap::BitmapBuilder;
26use risingwave_common::catalog::{CdcKeyComparison, ColumnDesc};
27use risingwave_common::row::RowExt;
28use risingwave_common::util::sort_util::OrderType;
29use risingwave_connector::parser::{
30    BigintUnsignedHandlingMode, ByteStreamSourceParser, DebeziumParser, DebeziumProps,
31    EncodingProperties, JsonProperties, ProtocolProperties, SourceStreamChunkBuilder,
32    SpecificParserConfig, TimeHandling, TimestampHandling, TimestamptzHandling,
33};
34use risingwave_connector::source::cdc::CdcScanOptions;
35use risingwave_connector::source::cdc::external::{
36    CdcOffset, ExternalCdcTableType, ExternalTableReaderImpl,
37};
38use risingwave_connector::source::{SourceColumnDesc, SourceContext, SourceCtrlOpts};
39use risingwave_pb::common::ThrottleType;
40use rw_futures_util::pausable;
41use thiserror_ext::AsReport;
42use tracing::Instrument;
43
44use crate::executor::backfill::cdc::state::CdcBackfillState;
45use crate::executor::backfill::cdc::upstream_table::external::ExternalStorageTable;
46use crate::executor::backfill::cdc::upstream_table::snapshot::{
47    SnapshotReadArgs, UpstreamTableRead, UpstreamTableReader,
48};
49use crate::executor::backfill::utils::{
50    cmp_pk_unsigned_aware, get_cdc_chunk_last_offset, get_new_pos, mapping_chunk, mapping_message,
51    mark_cdc_chunk,
52};
53use crate::executor::monitor::CdcBackfillMetrics;
54use crate::executor::prelude::*;
55use crate::executor::source::get_infinite_backoff_strategy;
56use crate::task::CreateMviewProgressReporter;
57
58/// `split_id`, `is_finished`, `row_count`, `cdc_offset` all occupy 1 column each.
59const METADATA_STATE_LEN: usize = 4;
60
61struct PkCompareInfo<'a> {
62    indices: &'a [usize],
63    order: &'a [OrderType],
64    needs_unsigned_i64_compare: &'a [bool],
65}
66
67fn can_poll_upstream_while_creating_reader(
68    current_pk_pos: Option<&OwnedRow>,
69    pk_comparisons_are_known: bool,
70) -> bool {
71    current_pk_pos.is_none() || pk_comparisons_are_known
72}
73
74// The TimestampHandling/TimestamptzHandling/TimeHandling parser's behavior depends on the debezium.time.precision.mode setting:
75// - If left unset, Debezium defaults to time.precision.mode=microseconds (per debezium.properties), and the parser uses Micro.
76// - If set to "connect", Debezium uses time.precision.mode=connect, and the parser uses Milli by design.
77// - If set to any other value, Debezium applies that specific value, and the default parser is used to maintain backward compatibility.
78pub(crate) fn get_cdc_json_parse_handling_from_properties(
79    properties: &BTreeMap<String, String>,
80) -> (
81    Option<TimestampHandling>,
82    Option<TimestamptzHandling>,
83    Option<TimeHandling>,
84    Option<BigintUnsignedHandlingMode>,
85) {
86    let (timestamp_handling, timestamptz_handling, time_handling) = match properties
87        .get("debezium.time.precision.mode")
88    {
89        None => (
90            Some(TimestampHandling::Micro),
91            Some(TimestamptzHandling::Micro),
92            Some(TimeHandling::Micro),
93        ),
94        Some(m) if m == "connect" => (
95            Some(TimestampHandling::Milli),
96            Some(TimestamptzHandling::Milli),
97            Some(TimeHandling::Milli),
98        ),
99        Some(other) => {
100            // backward compatibility.
101            tracing::warn!(
102                "Unsupported debezium.time.precision.mode = {other}, fall back to default parser."
103            );
104            (None, None, None)
105        }
106    };
107    let bigint_unsigned_handling = properties
108        .get("debezium.bigint.unsigned.handling.mode")
109        .is_some_and(|v| v == "precise")
110        .then_some(BigintUnsignedHandlingMode::Precise);
111    (
112        timestamp_handling,
113        timestamptz_handling,
114        time_handling,
115        bigint_unsigned_handling,
116    )
117}
118
119pub struct CdcBackfillExecutor<S: StateStore> {
120    actor_ctx: ActorContextRef,
121
122    /// The external table to be backfilled
123    external_table: ExternalStorageTable,
124
125    /// Upstream changelog stream which may contain metadata columns, e.g. `_rw_offset`
126    upstream: Executor,
127
128    /// The column indices need to be forwarded to the downstream from the upstream and table scan.
129    output_indices: Vec<usize>,
130
131    /// The schema of output chunk, including additional columns if any
132    output_columns: Vec<ColumnDesc>,
133
134    /// State table of the `CdcBackfill` executor
135    state_impl: CdcBackfillState<S>,
136
137    // TODO: introduce a CdcBackfillProgress to report finish to Meta
138    // This object is just a stub right now
139    progress: Option<CreateMviewProgressReporter>,
140
141    metrics: CdcBackfillMetrics,
142
143    /// Rate limit in rows/s.
144    rate_limit_rps: Option<u32>,
145
146    options: CdcScanOptions,
147
148    properties: BTreeMap<String, String>,
149}
150
151impl<S: StateStore> CdcBackfillExecutor<S> {
152    #[expect(clippy::too_many_arguments)]
153    pub fn new(
154        actor_ctx: ActorContextRef,
155        external_table: ExternalStorageTable,
156        upstream: Executor,
157        output_indices: Vec<usize>,
158        output_columns: Vec<ColumnDesc>,
159        progress: Option<CreateMviewProgressReporter>,
160        metrics: Arc<StreamingMetrics>,
161        state_table: StateTable<S>,
162        rate_limit_rps: Option<u32>,
163        options: CdcScanOptions,
164        properties: BTreeMap<String, String>,
165    ) -> Self {
166        let pk_indices = external_table.pk_indices();
167        let upstream_table_id = external_table.table_id();
168        let state_impl = CdcBackfillState::new(
169            upstream_table_id,
170            state_table,
171            pk_indices.len() + METADATA_STATE_LEN,
172        );
173
174        let metrics = metrics.new_cdc_backfill_metrics(external_table.table_id(), actor_ctx.id);
175        Self {
176            actor_ctx,
177            external_table,
178            upstream,
179            output_indices,
180            output_columns,
181            state_impl,
182            progress,
183            metrics,
184            rate_limit_rps,
185            options,
186            properties,
187        }
188    }
189
190    fn report_metrics(
191        metrics: &CdcBackfillMetrics,
192        snapshot_processed_row_count: u64,
193        upstream_processed_row_count: u64,
194    ) {
195        metrics
196            .cdc_backfill_snapshot_read_row_count
197            .inc_by(snapshot_processed_row_count);
198
199        metrics
200            .cdc_backfill_upstream_output_row_count
201            .inc_by(upstream_processed_row_count);
202    }
203
204    fn consume_upstream_chunk_buffer(
205        offset_parse_func: &risingwave_connector::source::cdc::external::CdcOffsetParseFunc,
206        upstream_chunk_buffer: &mut Vec<StreamChunk>,
207        current_pk_pos: Option<&OwnedRow>,
208        pk_compare: PkCompareInfo<'_>,
209        last_binlog_offset: &Option<CdcOffset>,
210        output_indices: &[usize],
211    ) -> StreamExecutorResult<(Vec<StreamChunk>, u64, Option<CdcOffset>)> {
212        let Some(current_pos) = current_pk_pos else {
213            return Ok((vec![], 0, None));
214        };
215
216        let buffered_chunks = std::mem::take(upstream_chunk_buffer);
217        let mut emitted_chunks = Vec::with_capacity(buffered_chunks.len());
218        let mut upstream_processed_row_count = 0;
219        let mut consumed_binlog_offset = None;
220        let mut retained_chunks = Vec::with_capacity(buffered_chunks.len());
221        let mut can_advance_consumed_offset = true;
222
223        for chunk in buffered_chunks {
224            let mut emitted_vis = BitmapBuilder::zeroed(chunk.capacity());
225            let mut retained_vis = BitmapBuilder::zeroed(chunk.capacity());
226            for idx in 0..chunk.capacity() {
227                let (op, row, visible) = chunk.row_at(idx);
228                if !visible {
229                    // Buffered chunks may carry sparse visibility from previous filtering rounds.
230                    continue;
231                }
232                let event_offset = (*offset_parse_func)(
233                    row.iter()
234                        .last()
235                        .flatten()
236                        .expect("cdc offset must exist")
237                        .into_utf8(),
238                )?;
239                let in_binlog_range = last_binlog_offset
240                    .as_ref()
241                    .is_none_or(|binlog_low| *binlog_low <= event_offset);
242
243                let row_pk = row.project(pk_compare.indices);
244                let reached_current_pos = cmp_pk_unsigned_aware(
245                    row_pk.iter(),
246                    current_pos.iter(),
247                    pk_compare.order,
248                    pk_compare.needs_unsigned_i64_compare,
249                )
250                .is_le();
251                if !in_binlog_range {
252                    continue;
253                }
254                let should_emit = reached_current_pos;
255
256                match op {
257                    Op::Insert | Op::Delete => {
258                        if should_emit {
259                            emitted_vis.set(idx, true);
260                            if can_advance_consumed_offset {
261                                consumed_binlog_offset = Some(event_offset);
262                            }
263                            upstream_processed_row_count += 1;
264                        } else {
265                            retained_vis.set(idx, true);
266                            can_advance_consumed_offset = false;
267                        }
268                    }
269                    Op::UpdateDelete | Op::UpdateInsert => {
270                        unreachable!("CDC buffered chunks should not contain update pairs")
271                    }
272                }
273            }
274
275            let emitted_vis = emitted_vis.finish();
276            if emitted_vis.count_ones() > 0 {
277                emitted_chunks.push(mapping_chunk(
278                    chunk.clone_with_vis(emitted_vis),
279                    output_indices,
280                ));
281            }
282
283            let retained_vis = retained_vis.finish();
284            if retained_vis.count_ones() > 0 {
285                retained_chunks.push(chunk.clone_with_vis(retained_vis));
286            }
287        }
288
289        *upstream_chunk_buffer = retained_chunks;
290
291        Ok((
292            emitted_chunks,
293            upstream_processed_row_count,
294            consumed_binlog_offset,
295        ))
296    }
297
298    fn filter_recovery_chunk(
299        offset_parse_func: &risingwave_connector::source::cdc::external::CdcOffsetParseFunc,
300        chunk: StreamChunk,
301        current_pos: &OwnedRow,
302        pk_compare: PkCompareInfo<'_>,
303        last_binlog_offset: &Option<CdcOffset>,
304        output_indices: &[usize],
305    ) -> StreamExecutorResult<(Option<StreamChunk>, Option<CdcOffset>)> {
306        let chunk_offset = get_cdc_chunk_last_offset(offset_parse_func, &chunk)?;
307        let chunk = mapping_chunk(
308            mark_cdc_chunk(
309                offset_parse_func,
310                chunk,
311                current_pos,
312                pk_compare.indices,
313                pk_compare.order,
314                pk_compare.needs_unsigned_i64_compare,
315                last_binlog_offset.clone(),
316            )?,
317            output_indices,
318        );
319        let chunk = (chunk.cardinality() > 0).then_some(chunk);
320        let consumed_offset = chunk_offset.filter(|chunk_offset| {
321            last_binlog_offset
322                .as_ref()
323                .is_none_or(|last| *last < *chunk_offset)
324        });
325        Ok((chunk, consumed_offset))
326    }
327
328    #[try_stream(ok = Message, error = StreamExecutorError)]
329    async fn execute_inner(mut self) {
330        // The indices to primary key columns
331        let pk_indices = self.external_table.pk_indices().to_vec();
332        let pk_order = self.external_table.pk_order_types().to_vec();
333        let mut pk_needs_unsigned_i64_compare =
334            self.external_table.pk_comparisons().map(|comparisons| {
335                comparisons
336                    .iter()
337                    .map(|comparison| *comparison == CdcKeyComparison::UnsignedInt64)
338                    .collect_vec()
339            });
340
341        let table_id = self.external_table.table_id();
342        let upstream_table_name = self.external_table.qualified_table_name();
343        let schema_table_name = self.external_table.schema_table_name().clone();
344        let external_database_name = self.external_table.database_name().to_owned();
345
346        let additional_columns = self
347            .output_columns
348            .iter()
349            .filter(|col| col.additional_column.column_type.is_some())
350            .cloned()
351            .collect_vec();
352
353        let mut upstream = self.upstream.execute();
354
355        // Current position of the upstream_table storage primary key.
356        // `None` means it starts from the beginning.
357        let mut current_pk_pos: Option<OwnedRow>;
358
359        // Poll the upstream to get the first barrier.
360        let first_barrier = expect_first_barrier(&mut upstream).await?;
361
362        let mut is_snapshot_paused = first_barrier.is_pause_on_startup();
363        let first_barrier_epoch = first_barrier.epoch;
364        // The first barrier message should be propagated.
365        yield Message::Barrier(first_barrier);
366
367        // Check whether this parallelism has been assigned splits,
368        // if not, we should bypass the backfill directly.
369        let mut state_impl = self.state_impl;
370
371        state_impl.init_epoch(first_barrier_epoch).await?;
372
373        // restore backfill state
374        let state = state_impl.restore_state().await?;
375        current_pk_pos = state.current_pk_pos.clone();
376
377        let need_backfill = !self.options.disable_backfill && !state.is_finished;
378
379        // Keep track of rows from the snapshot.
380        let mut total_snapshot_row_count = state.row_count as u64;
381
382        let (timestamp_handling, timestamptz_handling, time_handling, bigint_unsigned_handling) =
383            get_cdc_json_parse_handling_from_properties(&self.properties);
384        // Only postgres-cdc connector may trigger TOAST.
385        let handle_toast_columns: bool =
386            self.external_table.table_type() == &ExternalCdcTableType::Postgres;
387        // Make sure to use mapping_message after transform_upstream.
388        let mut upstream = transform_upstream(
389            upstream,
390            self.output_columns.clone(),
391            timestamp_handling,
392            timestamptz_handling,
393            time_handling,
394            bigint_unsigned_handling,
395            handle_toast_columns,
396        )
397        .boxed()
398        .peekable();
399        let mut last_binlog_offset = state.last_cdc_offset.clone();
400
401        // CDC Backfill Algorithm:
402        //
403        // When the first barrier comes from upstream:
404        //  - read the current binlog offset as `binlog_low`
405        //  - start a snapshot read upon upstream table and iterate over the snapshot read stream
406        //  - buffer the changelog event from upstream
407        //
408        // When a new barrier comes from upstream:
409        //  - read the current binlog offset as `binlog_high`
410        //  - for each row of the upstream change log, forward it to downstream if it in the range
411        //    of [binlog_low, binlog_high] and its pk <= `current_pos`, otherwise keep buffering it
412        //    until `current_pos` catches up
413        //  - reconstruct the whole backfill stream with upstream changelog and a new table snapshot
414        //
415        // When a chunk comes from snapshot, we forward it to the downstream and raise
416        // `current_pos`.
417        // When we reach the end of the snapshot read stream, it means backfill has been
418        // finished.
419        //
420        // Once the backfill loop ends, we forward the upstream directly to the downstream.
421        if need_backfill {
422            // After init the state table and forward the initial barrier to downstream,
423            // we now try to create the table reader with retry.
424            //
425            // A fresh backfill can ignore CDC events here because its snapshot starts from the
426            // beginning. Recovery preserves events for the already-scanned prefix while reader
427            // creation is retried. Legacy MySQL graphs with Int64 PK columns must first recover
428            // comparison metadata from the reader so events use unsigned-aware ordering if needed.
429            let offset_parse_func = self.external_table.table_type().get_cdc_offset_parser()?;
430            let mut table_reader = None;
431            let external_table = self.external_table.clone();
432            let actor_id = self.actor_ctx.id;
433            let fragment_id = self.actor_ctx.fragment_id;
434            let mut future = Box::pin(async move {
435                let backoff = get_infinite_backoff_strategy();
436                tokio_retry::Retry::spawn(backoff, || async {
437                    match external_table.create_table_reader().await {
438                        Ok(reader) => Ok(reader),
439                        Err(e) => {
440                            tracing::warn!(error = %e.as_report(), actor_id = %actor_id, fragment_id = %fragment_id, "failed to create cdc table reader, retrying...");
441                            Err(e)
442                        }
443                    }
444                })
445                .instrument(tracing::info_span!("create_cdc_table_reader_with_retry"))
446                .await
447                .expect("Retry create cdc table reader until success.")
448            });
449            if can_poll_upstream_while_creating_reader(
450                current_pk_pos.as_ref(),
451                pk_needs_unsigned_i64_compare.is_some(),
452            ) {
453                while let Some(msg) =
454                    build_reader_and_poll_upstream(&mut upstream, &mut table_reader, &mut future)
455                        .await?
456                {
457                    match msg {
458                        Message::Barrier(barrier) => {
459                            state_impl.commit_state(barrier.epoch).await?;
460                            yield Message::Barrier(barrier);
461                        }
462                        Message::Chunk(chunk) => {
463                            if let Some(current_pos) = current_pk_pos.as_ref() {
464                                let pk_needs_unsigned_i64_compare =
465                                    pk_needs_unsigned_i64_compare.as_deref().expect(
466                                        "recovery may only poll upstream with known PK comparisons",
467                                    );
468                                let (chunk, consumed_offset) = Self::filter_recovery_chunk(
469                                    &offset_parse_func,
470                                    chunk,
471                                    current_pos,
472                                    PkCompareInfo {
473                                        indices: &pk_indices,
474                                        order: &pk_order,
475                                        needs_unsigned_i64_compare: pk_needs_unsigned_i64_compare,
476                                    },
477                                    &last_binlog_offset,
478                                    &self.output_indices,
479                                )?;
480                                if let Some(chunk) = chunk {
481                                    Self::report_metrics(
482                                        &self.metrics,
483                                        0,
484                                        chunk.cardinality() as u64,
485                                    );
486                                    yield Message::Chunk(chunk);
487                                }
488                                if let Some(consumed_offset) = consumed_offset {
489                                    last_binlog_offset = Some(consumed_offset);
490                                }
491                            } else {
492                                // The snapshot starts from the beginning and covers these changes.
493                            }
494                        }
495                        Message::Watermark(_) => {
496                            // ignore watermark
497                        }
498                    }
499                }
500            } else {
501                tracing::warn!(
502                    %table_id,
503                    upstream_table_name,
504                    "waiting for the CDC table reader to recover legacy MySQL primary-key comparison metadata; checkpoint progress is paused"
505                );
506                table_reader = Some(future.as_mut().await);
507            }
508            let table_reader = table_reader.expect("table reader must be created");
509            if pk_needs_unsigned_i64_compare.is_none() {
510                let pk_names = pk_indices
511                    .iter()
512                    .map(|&idx| self.external_table.schema().fields[idx].name.clone())
513                    .collect_vec();
514                let comparisons = table_reader.pk_column_comparisons(&pk_names)?;
515                assert_eq!(comparisons.len(), pk_indices.len());
516                pk_needs_unsigned_i64_compare = Some(
517                    comparisons
518                        .into_iter()
519                        .map(|comparison| comparison == CdcKeyComparison::UnsignedInt64)
520                        .collect_vec(),
521                );
522            }
523            let pk_needs_unsigned_i64_compare = pk_needs_unsigned_i64_compare
524                .expect("PK comparison metadata must be resolved after reader creation");
525            tracing::info!(
526                %table_id,
527                upstream_table_name,
528                "table reader created successfully"
529            );
530
531            let upstream_table_reader =
532                UpstreamTableReader::new(self.external_table.clone(), table_reader);
533
534            if last_binlog_offset.is_none() {
535                // Limit concurrent CDC connections globally to 10 using a semaphore.
536                static CDC_CONN_SEMAPHORE: tokio::sync::Semaphore =
537                    tokio::sync::Semaphore::const_new(10);
538
539                let _permit = CDC_CONN_SEMAPHORE.acquire().await.unwrap();
540                last_binlog_offset = upstream_table_reader.current_cdc_offset().await?;
541            }
542
543            let mut consumed_binlog_offset: Option<CdcOffset> = None;
544
545            tracing::info!(
546                %table_id,
547                upstream_table_name,
548                initial_binlog_offset = ?last_binlog_offset,
549                ?current_pk_pos,
550                is_finished = state.is_finished,
551                is_snapshot_paused,
552                snapshot_row_count = total_snapshot_row_count,
553                rate_limit = self.rate_limit_rps,
554                disable_backfill = self.options.disable_backfill,
555                snapshot_barrier_interval = self.options.snapshot_barrier_interval,
556                snapshot_batch_size = self.options.snapshot_batch_size,
557                "start cdc backfill",
558            );
559
560            // drive the upstream changelog first to ensure we can receive timely changelog event,
561            // otherwise the upstream changelog may be blocked by the snapshot read stream
562            let _ = Pin::new(&mut upstream).peek().await;
563
564            // wait for a barrier to make sure the backfill starts after upstream source
565            #[for_await]
566            for msg in upstream.by_ref() {
567                match msg? {
568                    Message::Barrier(barrier) => {
569                        match barrier.mutation.as_deref() {
570                            Some(crate::executor::Mutation::Pause) => {
571                                is_snapshot_paused = true;
572                                tracing::info!(
573                                    %table_id,
574                                    upstream_table_name,
575                                    "snapshot is paused by barrier"
576                                );
577                            }
578                            Some(crate::executor::Mutation::Resume) => {
579                                is_snapshot_paused = false;
580                                tracing::info!(
581                                    %table_id,
582                                    upstream_table_name,
583                                    "snapshot is resumed by barrier"
584                                );
585                            }
586                            _ => {
587                                // ignore other mutations
588                            }
589                        }
590                        // Commit after all preceding recovery chunks have been emitted.
591                        state_impl.commit_state(barrier.epoch).await?;
592                        yield Message::Barrier(barrier);
593                        break;
594                    }
595                    Message::Chunk(chunk) => {
596                        if let Some(current_pos) = current_pk_pos.as_ref() {
597                            let (chunk, consumed_offset) = Self::filter_recovery_chunk(
598                                &offset_parse_func,
599                                chunk,
600                                current_pos,
601                                PkCompareInfo {
602                                    indices: &pk_indices,
603                                    order: &pk_order,
604                                    needs_unsigned_i64_compare: &pk_needs_unsigned_i64_compare,
605                                },
606                                &last_binlog_offset,
607                                &self.output_indices,
608                            )?;
609                            if let Some(chunk) = chunk {
610                                Self::report_metrics(&self.metrics, 0, chunk.cardinality() as u64);
611                                yield Message::Chunk(chunk);
612                            }
613                            if let Some(consumed_offset) = consumed_offset {
614                                last_binlog_offset = Some(consumed_offset);
615                            }
616                        } else if let Some(chunk_offset) =
617                            get_cdc_chunk_last_offset(&offset_parse_func, &chunk)?
618                            && last_binlog_offset
619                                .as_ref()
620                                .is_none_or(|last| *last < chunk_offset)
621                        {
622                            last_binlog_offset = Some(chunk_offset);
623                        }
624                    }
625                    Message::Watermark(_) => {
626                        // Ignore watermark
627                    }
628                }
629            }
630
631            tracing::info!(%table_id,
632                upstream_table_name,
633                initial_binlog_offset = ?last_binlog_offset,
634                ?current_pk_pos,
635                is_snapshot_paused,
636                "start cdc backfill loop");
637
638            // the buffer will be drained when a barrier comes
639            let mut upstream_chunk_buffer: Vec<StreamChunk> = vec![];
640
641            'backfill_loop: loop {
642                let mut should_bypass_snapshot_stream_patch =
643                    self.rate_limit_rps.is_some_and(|val| val == 0);
644                let left_upstream = upstream.by_ref().map(Either::Left);
645
646                let mut snapshot_read_row_cnt: usize = 0;
647                let read_args = SnapshotReadArgs::new(
648                    current_pk_pos.clone(),
649                    self.rate_limit_rps,
650                    pk_indices.clone(),
651                    additional_columns.clone(),
652                    schema_table_name.clone(),
653                    external_database_name.clone(),
654                );
655                let right_snapshot = pin!(
656                    upstream_table_reader
657                        .snapshot_read_full_table(read_args, self.options.snapshot_batch_size)
658                        .map(Either::Right)
659                );
660                let (right_snapshot, snapshot_valve) = pausable(right_snapshot);
661                if is_snapshot_paused {
662                    snapshot_valve.pause();
663                }
664
665                // Prefer to select upstream, so we can stop snapshot stream when barrier comes.
666                let mut backfill_stream =
667                    select_with_strategy(left_upstream, right_snapshot, |_: &mut ()| {
668                        stream::PollNext::Left
669                    });
670
671                let mut cur_barrier_snapshot_processed_rows: u64 = 0;
672                let mut cur_barrier_upstream_processed_rows: u64 = 0;
673                let mut barrier_count: u32 = 0;
674                let mut pending_barrier = None;
675
676                #[for_await]
677                for either in &mut backfill_stream {
678                    match either {
679                        // Upstream
680                        Either::Left(msg) => {
681                            match msg? {
682                                Message::Barrier(barrier) => {
683                                    // increase the barrier count and check whether need to start a new snapshot
684                                    barrier_count += 1;
685                                    let can_start_new_snapshot =
686                                        barrier_count == self.options.snapshot_barrier_interval;
687                                    let mut needs_rebuild_snapshot = false;
688
689                                    if let Some(mutation) = barrier.mutation.as_deref() {
690                                        use crate::executor::Mutation;
691                                        match mutation {
692                                            Mutation::Pause => {
693                                                is_snapshot_paused = true;
694                                                snapshot_valve.pause();
695                                            }
696                                            Mutation::Resume => {
697                                                is_snapshot_paused = false;
698                                                snapshot_valve.resume();
699                                            }
700                                            Mutation::Throttle(some) => {
701                                                if let Some(entry) =
702                                                    some.get(&self.actor_ctx.fragment_id)
703                                                    && entry.throttle_type()
704                                                        == ThrottleType::Backfill
705                                                    && entry.rate_limit != self.rate_limit_rps
706                                                {
707                                                    // If self.rate_limit_rps is 0, the snapshot stream does not establish a snapshot_read.
708                                                    // Consequently, the later snapshot stream patch is bypassed.
709                                                    should_bypass_snapshot_stream_patch = self
710                                                        .rate_limit_rps
711                                                        .is_some_and(|val| val == 0);
712                                                    self.rate_limit_rps = entry.rate_limit;
713                                                    needs_rebuild_snapshot = true;
714                                                }
715                                            }
716                                            mutation if mutation.is_stop(self.actor_ctx.id) => {
717                                                // the actor has been dropped, exit the backfill loop
718                                                tracing::info!(
719                                                    %table_id,
720                                                    upstream_table_name,
721                                                    "CdcBackfill has been dropped due to config change"
722                                                );
723                                                yield Message::Barrier(barrier);
724                                                break 'backfill_loop;
725                                            }
726                                            _ => (),
727                                        }
728                                    }
729
730                                    // when processing a barrier, check whether can start a new snapshot
731                                    // if the number of barriers reaches the snapshot interval
732                                    if can_start_new_snapshot || needs_rebuild_snapshot {
733                                        // staging the barrier
734                                        pending_barrier = Some(barrier);
735                                        tracing::debug!(
736                                            %table_id,
737                                            ?current_pk_pos,
738                                            ?snapshot_read_row_cnt,
739                                            "Prepare to start a new snapshot"
740                                        );
741                                        // Break the loop for consuming snapshot and prepare to start a new snapshot
742                                        break;
743                                    } else {
744                                        // Drain the in-memory buffer to ensure no data is lost during the recovery process.
745                                        let (
746                                            emitted_upstream_chunks,
747                                            consumed_upstream_row_count,
748                                            consumed_binlog_offset,
749                                        ) = Self::consume_upstream_chunk_buffer(
750                                            &offset_parse_func,
751                                            &mut upstream_chunk_buffer,
752                                            current_pk_pos.as_ref(),
753                                            PkCompareInfo {
754                                                indices: &pk_indices,
755                                                order: &pk_order,
756                                                needs_unsigned_i64_compare:
757                                                    &pk_needs_unsigned_i64_compare,
758                                            },
759                                            &last_binlog_offset,
760                                            &self.output_indices,
761                                        )?;
762                                        cur_barrier_upstream_processed_rows +=
763                                            consumed_upstream_row_count;
764                                        if let Some(consumed_binlog_offset) = consumed_binlog_offset
765                                        {
766                                            last_binlog_offset = Some(consumed_binlog_offset);
767                                        }
768                                        for chunk in emitted_upstream_chunks {
769                                            yield Message::Chunk(chunk);
770                                        }
771
772                                        Self::report_metrics(
773                                            &self.metrics,
774                                            cur_barrier_snapshot_processed_rows,
775                                            cur_barrier_upstream_processed_rows,
776                                        );
777
778                                        // update and persist current backfill progress
779                                        state_impl
780                                            .mutate_state(
781                                                current_pk_pos.clone(),
782                                                last_binlog_offset.clone(),
783                                                total_snapshot_row_count,
784                                                false,
785                                            )
786                                            .await?;
787
788                                        state_impl.commit_state(barrier.epoch).await?;
789
790                                        // emit barrier and continue consume the backfill stream
791                                        yield Message::Barrier(barrier);
792                                    }
793                                }
794                                Message::Chunk(chunk) => {
795                                    // skip empty upstream chunk
796                                    if chunk.cardinality() == 0 {
797                                        continue;
798                                    }
799
800                                    let chunk_binlog_offset =
801                                        get_cdc_chunk_last_offset(&offset_parse_func, &chunk)?;
802
803                                    tracing::trace!(
804                                        "recv changelog chunk: chunk_offset {:?}, capactiy {}",
805                                        chunk_binlog_offset,
806                                        chunk.capacity()
807                                    );
808
809                                    // Since we don't need changelog before the
810                                    // `last_binlog_offset`, skip the chunk that *only* contains
811                                    // events before `last_binlog_offset`.
812                                    if let Some(last_binlog_offset) = last_binlog_offset.as_ref()
813                                        && let Some(chunk_offset) = chunk_binlog_offset
814                                        && chunk_offset < *last_binlog_offset
815                                    {
816                                        tracing::trace!(
817                                            "skip changelog chunk: chunk_offset {:?}, capacity {}",
818                                            chunk_offset,
819                                            chunk.capacity()
820                                        );
821                                        continue;
822                                    }
823                                    // Buffer the upstream chunk.
824                                    upstream_chunk_buffer.push(chunk.compact_vis());
825                                }
826                                Message::Watermark(_) => {
827                                    // Ignore watermark during backfill.
828                                }
829                            }
830                        }
831                        // Snapshot read
832                        Either::Right(msg) => {
833                            match msg? {
834                                None => {
835                                    tracing::info!(
836                                        %table_id,
837                                        ?last_binlog_offset,
838                                        ?current_pk_pos,
839                                        "snapshot read stream ends"
840                                    );
841                                    // If the snapshot read stream ends, it means all historical
842                                    // data has been loaded.
843                                    // We should not mark the chunk anymore,
844                                    // otherwise, we will ignore some rows in the buffer.
845                                    for chunk in upstream_chunk_buffer.drain(..) {
846                                        yield Message::Chunk(mapping_chunk(
847                                            chunk,
848                                            &self.output_indices,
849                                        ));
850                                    }
851
852                                    // backfill has finished, exit the backfill loop and persist the state when we recv a barrier
853                                    break 'backfill_loop;
854                                }
855                                Some(chunk) => {
856                                    // Raise the current position.
857                                    // As snapshot read streams are ordered by pk, so we can
858                                    // just use the last row to update `current_pos`.
859                                    current_pk_pos = Some(get_new_pos(&chunk, &pk_indices));
860
861                                    tracing::trace!(
862                                        "got a snapshot chunk: len {}, current_pk_pos {:?}",
863                                        chunk.cardinality(),
864                                        current_pk_pos
865                                    );
866                                    let chunk_cardinality = chunk.cardinality() as u64;
867                                    cur_barrier_snapshot_processed_rows += chunk_cardinality;
868                                    total_snapshot_row_count += chunk_cardinality;
869                                    yield Message::Chunk(mapping_chunk(
870                                        chunk,
871                                        &self.output_indices,
872                                    ));
873                                }
874                            }
875                        }
876                    }
877                }
878
879                assert!(pending_barrier.is_some(), "pending_barrier must exist");
880                let pending_barrier = pending_barrier.unwrap();
881
882                // The snapshot stream patch:
883                // Here we have to ensure the snapshot stream is consumed at least once,
884                // since the barrier event can kick in anytime.
885                // Otherwise, the result set of the new snapshot stream may become empty.
886                // It maybe a cancellation bug of the mysql driver.
887                let (_, mut snapshot_stream) = backfill_stream.into_inner();
888                // Resume the snapshot stream so that the snapshot stream patch won't block.
889                if !should_bypass_snapshot_stream_patch && is_snapshot_paused {
890                    snapshot_valve.resume();
891                }
892                if !should_bypass_snapshot_stream_patch
893                    && let Some(msg) = snapshot_stream
894                        .next()
895                        .instrument_await("consume_snapshot_stream_once")
896                        .await
897                {
898                    let Either::Right(msg) = msg else {
899                        bail!("BUG: snapshot_read contains upstream messages");
900                    };
901                    match msg? {
902                        None => {
903                            tracing::info!(
904                                %table_id,
905                                ?last_binlog_offset,
906                                ?current_pk_pos,
907                                "snapshot read stream ends in the force emit branch"
908                            );
909                            // End of the snapshot read stream.
910                            // Consume the buffered upstream chunk without filtering by `binlog_low`.
911                            for chunk in upstream_chunk_buffer.drain(..) {
912                                yield Message::Chunk(mapping_chunk(chunk, &self.output_indices));
913                            }
914
915                            // mark backfill has finished
916                            state_impl
917                                .mutate_state(
918                                    current_pk_pos.clone(),
919                                    last_binlog_offset.clone(),
920                                    total_snapshot_row_count,
921                                    true,
922                                )
923                                .await?;
924
925                            // commit state because we have received a barrier message
926                            state_impl.commit_state(pending_barrier.epoch).await?;
927                            yield Message::Barrier(pending_barrier);
928                            // end of backfill loop, since backfill has finished
929                            break 'backfill_loop;
930                        }
931                        Some(_) if is_snapshot_paused => {
932                            // Since the snapshot stream is paused, drop the chunk.
933                        }
934                        Some(chunk) => {
935                            // Raise the current pk position.
936                            current_pk_pos = Some(get_new_pos(&chunk, &pk_indices));
937
938                            let row_count = chunk.cardinality() as u64;
939                            cur_barrier_snapshot_processed_rows += row_count;
940                            total_snapshot_row_count += row_count;
941                            snapshot_read_row_cnt += row_count as usize;
942
943                            tracing::debug!(
944                                %table_id,
945                                ?current_pk_pos,
946                                ?snapshot_read_row_cnt,
947                                "force emit a snapshot chunk"
948                            );
949                            yield Message::Chunk(mapping_chunk(chunk, &self.output_indices));
950                        }
951                    }
952                }
953
954                // If the number of barriers reaches the snapshot interval,
955                // consume the buffered upstream chunks.
956                if let Some(current_pos) = &current_pk_pos {
957                    for chunk in upstream_chunk_buffer.drain(..) {
958                        cur_barrier_upstream_processed_rows += chunk.cardinality() as u64;
959
960                        // record the consumed binlog offset that will be
961                        // persisted later
962                        consumed_binlog_offset =
963                            get_cdc_chunk_last_offset(&offset_parse_func, &chunk)?;
964
965                        yield Message::Chunk(mapping_chunk(
966                            mark_cdc_chunk(
967                                &offset_parse_func,
968                                chunk,
969                                current_pos,
970                                &pk_indices,
971                                &pk_order,
972                                &pk_needs_unsigned_i64_compare,
973                                last_binlog_offset.clone(),
974                            )?,
975                            &self.output_indices,
976                        ));
977                    }
978                } else {
979                    // If no current_pos, means we did not process any snapshot yet.
980                    // we can just ignore the upstream buffer chunk in that case.
981                    upstream_chunk_buffer.clear();
982                }
983
984                // Update last seen binlog offset
985                if consumed_binlog_offset.is_some() {
986                    last_binlog_offset.clone_from(&consumed_binlog_offset);
987                }
988
989                Self::report_metrics(
990                    &self.metrics,
991                    cur_barrier_snapshot_processed_rows,
992                    cur_barrier_upstream_processed_rows,
993                );
994
995                // update and persist current backfill progress
996                state_impl
997                    .mutate_state(
998                        current_pk_pos.clone(),
999                        last_binlog_offset.clone(),
1000                        total_snapshot_row_count,
1001                        false,
1002                    )
1003                    .await?;
1004
1005                state_impl.commit_state(pending_barrier.epoch).await?;
1006                yield Message::Barrier(pending_barrier);
1007            }
1008            upstream_table_reader.disconnect().await?;
1009        } else if self.options.disable_backfill {
1010            // If backfill is disabled, we just mark the backfill as finished
1011            tracing::info!(
1012                %table_id,
1013                upstream_table_name,
1014                "CdcBackfill has been disabled"
1015            );
1016            state_impl
1017                .mutate_state(
1018                    current_pk_pos.clone(),
1019                    last_binlog_offset.clone(),
1020                    total_snapshot_row_count,
1021                    true,
1022                )
1023                .await?;
1024        }
1025
1026        tracing::info!(
1027            %table_id,
1028            upstream_table_name,
1029            "CdcBackfill has already finished and will forward messages directly to the downstream"
1030        );
1031
1032        // Wait for first barrier to come after backfill is finished.
1033        // So we can update our progress + persist the status.
1034        while let Some(Ok(msg)) = upstream.next().await {
1035            if let Some(msg) = mapping_message(msg, &self.output_indices) {
1036                // If not finished then we need to update state, otherwise no need.
1037                if let Message::Barrier(barrier) = &msg {
1038                    // finalized the backfill state
1039                    // TODO: unify `mutate_state` and `commit_state` into one method
1040                    state_impl
1041                        .mutate_state(
1042                            current_pk_pos.clone(),
1043                            last_binlog_offset.clone(),
1044                            total_snapshot_row_count,
1045                            true,
1046                        )
1047                        .await?;
1048                    state_impl.commit_state(barrier.epoch).await?;
1049
1050                    // mark progress as finished
1051                    if let Some(progress) = self.progress.as_mut() {
1052                        progress.finish(barrier.epoch, total_snapshot_row_count);
1053                    }
1054                    yield msg;
1055                    // break after the state have been saved
1056                    break;
1057                }
1058                yield msg;
1059            }
1060        }
1061
1062        // After backfill progress finished
1063        // we can forward messages directly to the downstream,
1064        // as backfill is finished.
1065        #[for_await]
1066        for msg in upstream {
1067            // upstream offsets will be removed from the message before forwarding to
1068            // downstream
1069            if let Some(msg) = mapping_message(msg?, &self.output_indices) {
1070                if let Message::Barrier(barrier) = &msg {
1071                    // commit state just to bump the epoch of state table
1072                    state_impl.commit_state(barrier.epoch).await?;
1073                }
1074                yield msg;
1075            }
1076        }
1077    }
1078}
1079
1080pub(crate) async fn build_reader_and_poll_upstream(
1081    upstream: &mut (impl Stream<Item = StreamExecutorResult<Message>> + Unpin),
1082    table_reader: &mut Option<ExternalTableReaderImpl>,
1083    future: &mut Pin<Box<impl Future<Output = ExternalTableReaderImpl>>>,
1084) -> StreamExecutorResult<Option<Message>> {
1085    if table_reader.is_some() {
1086        return Ok(None);
1087    }
1088    tokio::select! {
1089        biased;
1090        reader = &mut *future => {
1091            *table_reader = Some(reader);
1092            Ok(None)
1093        }
1094        msg = upstream.next() => {
1095            msg.transpose()
1096        }
1097    }
1098}
1099
1100#[try_stream(ok = Message, error = StreamExecutorError)]
1101pub async fn transform_upstream(
1102    upstream: BoxedMessageStream,
1103    output_columns: Vec<ColumnDesc>,
1104    timestamp_handling: Option<TimestampHandling>,
1105    timestamptz_handling: Option<TimestamptzHandling>,
1106    time_handling: Option<TimeHandling>,
1107    bigint_unsigned_handling: Option<BigintUnsignedHandlingMode>,
1108    handle_toast_columns: bool,
1109) {
1110    let props = SpecificParserConfig {
1111        encoding_config: EncodingProperties::Json(JsonProperties {
1112            use_schema_registry: false,
1113            timestamp_handling,
1114            timestamptz_handling,
1115            time_handling,
1116            bigint_unsigned_handling,
1117            handle_toast_columns,
1118        }),
1119        // the cdc message is generated internally so the key must exist.
1120        protocol_config: ProtocolProperties::Debezium(DebeziumProps::default()),
1121    };
1122
1123    // convert to source column desc to feed into parser
1124    let columns_with_meta = output_columns
1125        .iter()
1126        .map(SourceColumnDesc::from)
1127        .collect_vec();
1128    let mut parser = DebeziumParser::new(
1129        props,
1130        columns_with_meta.clone(),
1131        Arc::new(SourceContext::dummy()),
1132    )
1133    .await
1134    .map_err(StreamExecutorError::connector_error)?;
1135
1136    pin_mut!(upstream);
1137    #[for_await]
1138    for msg in upstream {
1139        let mut msg = msg?;
1140        if let Message::Chunk(chunk) = &mut msg {
1141            let parsed_chunk = parse_debezium_chunk(&mut parser, chunk).await?;
1142            let _ = std::mem::replace(chunk, parsed_chunk);
1143        }
1144        yield msg;
1145    }
1146}
1147
1148async fn parse_debezium_chunk(
1149    parser: &mut DebeziumParser,
1150    chunk: &StreamChunk,
1151) -> StreamExecutorResult<StreamChunk> {
1152    // here we transform the input chunk in `(payload varchar, _rw_offset varchar, _rw_table_name varchar)` schema
1153    // to chunk with downstream table schema `info.schema` of MergeNode contains the schema of the
1154    // table job with `_rw_offset` in the end
1155    // see `gen_create_table_plan_for_cdc_source` for details
1156
1157    // use `SourceStreamChunkBuilder` for convenience
1158    let mut builder = SourceStreamChunkBuilder::new(
1159        parser.columns().to_vec(),
1160        SourceCtrlOpts {
1161            chunk_size: chunk.capacity(),
1162            split_txn: false,
1163        },
1164    );
1165
1166    // The schema of input chunk `(payload varchar, _rw_offset varchar, _rw_table_name varchar, _row_id)`
1167    // We should use the debezium parser to parse the first column,
1168    // then chain the parsed row with `_rw_offset` row to get a new row.
1169    let payloads = chunk.data_chunk().project(&[0]);
1170    let offsets = chunk.data_chunk().project(&[1]).compact_vis();
1171
1172    // TODO: preserve the transaction semantics
1173    for payload in payloads.rows() {
1174        let ScalarRefImpl::Jsonb(jsonb_ref) = payload.datum_at(0).expect("payload must exist")
1175        else {
1176            panic!("payload must be jsonb");
1177        };
1178
1179        parser
1180            .parse_inner(
1181                None,
1182                Some(jsonb_ref.to_string().as_bytes().to_vec()),
1183                builder.row_writer(),
1184            )
1185            .await
1186            .unwrap();
1187    }
1188    builder.finish_current_chunk();
1189
1190    let parsed_chunk = {
1191        let mut iter = builder.consume_ready_chunks();
1192        assert_eq!(1, iter.len());
1193        iter.next().unwrap()
1194    };
1195    assert_eq!(parsed_chunk.capacity(), chunk.capacity()); // each payload is expected to generate one row
1196    let (ops, mut columns, vis) = parsed_chunk.into_inner();
1197    // note that `vis` is not necessarily the same as the original chunk's visibilities
1198
1199    // concat the rows in the parsed chunk with the `_rw_offset` column
1200    columns.extend(offsets.into_parts().0);
1201
1202    Ok(StreamChunk::from_parts(
1203        ops,
1204        DataChunk::from_parts(columns.into(), vis),
1205    ))
1206}
1207
1208impl<S: StateStore> Execute for CdcBackfillExecutor<S> {
1209    fn execute(self: Box<Self>) -> BoxedMessageStream {
1210        self.execute_inner().boxed()
1211    }
1212}
1213
1214#[cfg(test)]
1215mod tests {
1216    use std::collections::BTreeMap;
1217    use std::str::FromStr;
1218
1219    use futures::{StreamExt, pin_mut, stream};
1220    use risingwave_common::array::{Array, DataChunk, Op, StreamChunk};
1221    use risingwave_common::catalog::{
1222        CdcKeyComparison, ColumnDesc, ColumnId, Field, Schema, TableId,
1223    };
1224    use risingwave_common::row::{OwnedRow, Row};
1225    use risingwave_common::types::{DataType, Datum, JsonbVal, ScalarImpl};
1226    use risingwave_common::util::epoch::test_epoch;
1227    use risingwave_common::util::iter_util::ZipEqFast;
1228    use risingwave_common::util::sort_util::OrderType;
1229    use risingwave_connector::source::cdc::CdcScanOptions;
1230    use risingwave_connector::source::cdc::external::mock_external_table::MockExternalTableReader;
1231    use risingwave_connector::source::cdc::external::mysql::MySqlOffset;
1232    use risingwave_connector::source::cdc::external::{
1233        CdcOffset, ExternalCdcTableType, ExternalTableConfig, ExternalTableReaderImpl,
1234        SchemaTableName,
1235    };
1236    use risingwave_storage::memory::MemoryStateStore;
1237
1238    use super::{
1239        PkCompareInfo, build_reader_and_poll_upstream, can_poll_upstream_while_creating_reader,
1240    };
1241    use crate::common::table::test_utils::gen_pbtable;
1242    use crate::executor::backfill::cdc::cdc_backfill::transform_upstream;
1243    use crate::executor::backfill::cdc::state::CdcBackfillState;
1244    use crate::executor::monitor::StreamingMetrics;
1245    use crate::executor::prelude::StateTable;
1246    use crate::executor::source::default_source_internal_table;
1247    use crate::executor::test_utils::{MessageSender, MockSource};
1248    use crate::executor::{
1249        ActorContext, Barrier, CdcBackfillExecutor, ExternalStorageTable, Message,
1250    };
1251
1252    #[tokio::test]
1253    async fn test_transform_upstream_chunk() {
1254        let schema = Schema::new(vec![
1255            Field::unnamed(DataType::Jsonb),   // debezium json payload
1256            Field::unnamed(DataType::Varchar), // _rw_offset
1257            Field::unnamed(DataType::Varchar), // _rw_table_name
1258        ]);
1259        let stream_key = vec![1];
1260        let (mut tx, source) = MockSource::channel();
1261        let source = source.into_executor(schema.clone(), stream_key.clone());
1262        // let payload = r#"{"before": null,"after":{"O_ORDERKEY": 5, "O_CUSTKEY": 44485, "O_ORDERSTATUS": "F", "O_TOTALPRICE": "144659.20", "O_ORDERDATE": "1994-07-30" },"source":{"version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002", "ts_ms": 1695277757000, "snapshot": "last", "db": "mydb", "sequence": null, "table": "orders_new", "server_id": 0, "gtid": null, "file": "binlog.000008", "pos": 3693, "row": 0, "thread": null, "query": null},"op":"r","ts_ms":1695277757017,"transaction":null}"#.to_string();
1263        let payload = r#"{ "payload": { "before": null, "after": { "O_ORDERKEY": 5, "O_CUSTKEY": 44485, "O_ORDERSTATUS": "F", "O_TOTALPRICE": "144659.20", "O_ORDERDATE": "1994-07-30" }, "source": { "version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002", "ts_ms": 1695277757000, "snapshot": "last", "db": "mydb", "sequence": null, "table": "orders_new", "server_id": 0, "gtid": null, "file": "binlog.000008", "pos": 3693, "row": 0, "thread": null, "query": null }, "op": "r", "ts_ms": 1695277757017, "transaction": null } }"#;
1264
1265        let datums: Vec<Datum> = vec![
1266            Some(JsonbVal::from_str(payload).unwrap().into()),
1267            Some("file: 1.binlog, pos: 100".to_owned().into()),
1268            Some("mydb.orders".to_owned().into()),
1269        ];
1270
1271        let mut builders = schema.create_array_builders(8);
1272        for (builder, datum) in builders.iter_mut().zip_eq_fast(datums.iter()) {
1273            builder.append(datum.clone());
1274        }
1275        let columns = builders
1276            .into_iter()
1277            .map(|builder| builder.finish().into())
1278            .collect();
1279
1280        // one row chunk
1281        let chunk = StreamChunk::from_parts(vec![Op::Insert], DataChunk::new(columns, 1));
1282
1283        tx.push_chunk(chunk);
1284        let upstream = Box::new(source).execute();
1285
1286        // schema to the debezium parser
1287        let columns = vec![
1288            ColumnDesc::named("O_ORDERKEY", ColumnId::new(1), DataType::Int64),
1289            ColumnDesc::named("O_CUSTKEY", ColumnId::new(2), DataType::Int64),
1290            ColumnDesc::named("O_ORDERSTATUS", ColumnId::new(3), DataType::Varchar),
1291            ColumnDesc::named("O_TOTALPRICE", ColumnId::new(4), DataType::Decimal),
1292            ColumnDesc::named("O_ORDERDATE", ColumnId::new(5), DataType::Date),
1293            ColumnDesc::named("commit_ts", ColumnId::new(6), DataType::Timestamptz),
1294        ];
1295
1296        let parsed_stream = transform_upstream(upstream, columns, None, None, None, None, false);
1297        pin_mut!(parsed_stream);
1298        let message = parsed_stream
1299            .next()
1300            .await
1301            .expect("transform stream should yield the input chunk")
1302            .expect("transforming the CDC chunk should succeed");
1303        let Message::Chunk(chunk) = message else {
1304            panic!("expected a transformed chunk");
1305        };
1306        assert_eq!(
1307            chunk
1308                .rows()
1309                .map(|(op, row)| (op, row.to_owned_row()))
1310                .collect::<Vec<_>>(),
1311            vec![(
1312                Op::Insert,
1313                OwnedRow::new(vec![
1314                    Some(ScalarImpl::Int64(5)),
1315                    Some(ScalarImpl::Int64(44485)),
1316                    Some(ScalarImpl::Utf8("F".into())),
1317                    Some(ScalarImpl::Decimal("144659.20".parse().unwrap())),
1318                    Some(ScalarImpl::Date("1994-07-30".parse().unwrap())),
1319                    None,
1320                    Some(ScalarImpl::Utf8("file: 1.binlog, pos: 100".into())),
1321                ]),
1322            )]
1323        );
1324    }
1325
1326    #[tokio::test]
1327    async fn test_build_reader_and_poll_upstream() {
1328        let actor_context = ActorContext::for_test(1);
1329        let external_storage_table = ExternalStorageTable::for_test_undefined();
1330        let schema = Schema::new(vec![
1331            Field::unnamed(DataType::Jsonb),   // debezium json payload
1332            Field::unnamed(DataType::Varchar), // _rw_offset
1333            Field::unnamed(DataType::Varchar), // _rw_table_name
1334        ]);
1335        let stream_key = vec![1];
1336        let (mut tx, source) = MockSource::channel();
1337        let source = source.into_executor(schema.clone(), stream_key.clone());
1338        let output_indices = vec![1, 0, 4]; //reorder
1339        let output_columns = vec![
1340            ColumnDesc::named("O_ORDERKEY", ColumnId::new(1), DataType::Int64),
1341            ColumnDesc::named("O_CUSTKEY", ColumnId::new(2), DataType::Int64),
1342            ColumnDesc::named("O_ORDERSTATUS", ColumnId::new(3), DataType::Varchar),
1343            ColumnDesc::named("O_TOTALPRICE", ColumnId::new(4), DataType::Decimal),
1344            ColumnDesc::named("O_DUMMY", ColumnId::new(5), DataType::Int64),
1345            ColumnDesc::named("commit_ts", ColumnId::new(6), DataType::Timestamptz),
1346        ];
1347        let store = MemoryStateStore::new();
1348        let state_table =
1349            StateTable::from_table_catalog(&default_source_internal_table(0x2333), store, None)
1350                .await;
1351        let cdc = CdcBackfillExecutor::new(
1352            actor_context,
1353            external_storage_table,
1354            source,
1355            output_indices,
1356            output_columns,
1357            None,
1358            StreamingMetrics::unused().into(),
1359            state_table,
1360            None,
1361            CdcScanOptions {
1362                // We want to mark backfill as finished. However it's not straightforward to do so.
1363                // Here we disable_backfill instead.
1364                disable_backfill: true,
1365                ..CdcScanOptions::default()
1366            },
1367            BTreeMap::default(),
1368        );
1369        // cdc.state_impl.init_epoch(EpochPair::new(test_epoch(4), test_epoch(3))).await.unwrap();
1370        // cdc.state_impl.mutate_state(None, None, 0, true).await.unwrap();
1371        // cdc.state_impl.commit_state(EpochPair::new(test_epoch(5), test_epoch(4))).await.unwrap();
1372        let s = cdc.execute_inner();
1373        pin_mut!(s);
1374
1375        // send first barrier
1376        tx.send_barrier(Barrier::new_test_barrier(test_epoch(8)));
1377        // send chunk
1378        {
1379            let payload = r#"{ "payload": { "before": null, "after": { "O_ORDERKEY": 5, "O_CUSTKEY": 44485, "O_ORDERSTATUS": "F", "O_TOTALPRICE": "144659.20", "O_DUMMY": 100 }, "source": { "version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002", "ts_ms": 1695277757000, "snapshot": "last", "db": "mydb", "sequence": null, "table": "orders_new", "server_id": 0, "gtid": null, "file": "binlog.000008", "pos": 3693, "row": 0, "thread": null, "query": null }, "op": "r", "ts_ms": 1695277757017, "transaction": null } }"#;
1380            let datums: Vec<Datum> = vec![
1381                Some(JsonbVal::from_str(payload).unwrap().into()),
1382                Some("file: 1.binlog, pos: 100".to_owned().into()),
1383                Some("mydb.orders".to_owned().into()),
1384            ];
1385            let mut builders = schema.create_array_builders(8);
1386            for (builder, datum) in builders.iter_mut().zip_eq_fast(datums.iter()) {
1387                builder.append(datum.clone());
1388            }
1389            let columns = builders
1390                .into_iter()
1391                .map(|builder| builder.finish().into())
1392                .collect();
1393            // one row chunk
1394            let chunk = StreamChunk::from_parts(vec![Op::Insert], DataChunk::new(columns, 1));
1395
1396            tx.push_chunk(chunk);
1397        }
1398        let _first_barrier = s.next().await.unwrap();
1399        let upstream_change_log = s.next().await.unwrap().unwrap();
1400        let Message::Chunk(chunk) = upstream_change_log else {
1401            panic!("expect chunk");
1402        };
1403        assert_eq!(chunk.columns().len(), 3);
1404        assert_eq!(chunk.rows().count(), 1);
1405        assert_eq!(
1406            chunk.columns()[0].as_int64().iter().collect::<Vec<_>>(),
1407            vec![Some(44485)]
1408        );
1409        assert_eq!(
1410            chunk.columns()[1].as_int64().iter().collect::<Vec<_>>(),
1411            vec![Some(5)]
1412        );
1413        assert_eq!(
1414            chunk.columns()[2].as_int64().iter().collect::<Vec<_>>(),
1415            vec![Some(100)]
1416        );
1417    }
1418
1419    #[test]
1420    fn test_reader_startup_polling_selection() {
1421        let recovered_position = OwnedRow::new(vec![Some(ScalarImpl::Int64(1))]);
1422
1423        assert!(!can_poll_upstream_while_creating_reader(
1424            Some(&recovered_position),
1425            false,
1426        ));
1427        assert!(can_poll_upstream_while_creating_reader(None, false));
1428        assert!(can_poll_upstream_while_creating_reader(
1429            Some(&recovered_position),
1430            true,
1431        ));
1432    }
1433
1434    #[tokio::test]
1435    async fn test_disabled_cdc_backfill_finishes_without_table_reader() {
1436        let actor_context = ActorContext::for_test(1);
1437        let external_storage_table = ExternalStorageTable::new(
1438            TableId::new(1234),
1439            SchemaTableName {
1440                schema_name: "public".to_owned(),
1441                table_name: "orders".to_owned(),
1442            },
1443            "db".to_owned(),
1444            ExternalTableConfig::default(),
1445            ExternalCdcTableType::Undefined,
1446            Schema::new(vec![Field::with_name(DataType::Int64, "id")]),
1447            vec![OrderType::ascending()],
1448            Some(vec![CdcKeyComparison::Native]),
1449            vec![0],
1450        );
1451        let schema = Schema::new(vec![
1452            Field::unnamed(DataType::Jsonb),
1453            Field::unnamed(DataType::Varchar),
1454            Field::unnamed(DataType::Varchar),
1455        ]);
1456        let stream_key = vec![1];
1457        let (tx, source) = MockSource::channel();
1458        let source = source.into_executor(schema, stream_key);
1459        let output_columns = vec![ColumnDesc::named("id", ColumnId::new(1), DataType::Int64)];
1460        let memory_state_store = MemoryStateStore::new();
1461        let state_table = create_cdc_state_table(memory_state_store.clone()).await;
1462        let cdc = CdcBackfillExecutor::new(
1463            actor_context,
1464            external_storage_table,
1465            source,
1466            vec![0],
1467            output_columns,
1468            None,
1469            StreamingMetrics::unused().into(),
1470            state_table,
1471            None,
1472            CdcScanOptions {
1473                disable_backfill: true,
1474                ..CdcScanOptions::default()
1475            },
1476            BTreeMap::default(),
1477        );
1478        let executor = cdc.execute_inner();
1479        pin_mut!(executor);
1480
1481        tx.send_barrier(Barrier::new_test_barrier(test_epoch(1)));
1482        assert!(matches!(
1483            executor.next().await.unwrap().unwrap(),
1484            Message::Barrier(_)
1485        ));
1486
1487        tx.send_barrier(Barrier::new_test_barrier(test_epoch(2)));
1488        assert!(matches!(
1489            executor.next().await.unwrap().unwrap(),
1490            Message::Barrier(_)
1491        ));
1492
1493        let mut restored_state = CdcBackfillState::new(
1494            TableId::new(1234),
1495            create_cdc_state_table(memory_state_store).await,
1496            5,
1497        );
1498        restored_state
1499            .init_epoch(Barrier::new_test_barrier(test_epoch(2)).epoch)
1500            .await
1501            .unwrap();
1502        let state = restored_state.restore_state().await.unwrap();
1503        assert!(state.is_finished);
1504        assert_eq!(state.last_cdc_offset, None);
1505    }
1506
1507    #[tokio::test]
1508    async fn test_unfinished_cdc_backfill_rejects_null_offset_on_restore() {
1509        let memory_state_store = MemoryStateStore::new();
1510        let mut state_writer = CdcBackfillState::new(
1511            TableId::new(1234),
1512            create_cdc_state_table(memory_state_store.clone()).await,
1513            5,
1514        );
1515        state_writer
1516            .init_epoch(Barrier::new_test_barrier(test_epoch(1)).epoch)
1517            .await
1518            .unwrap();
1519        state_writer
1520            .mutate_state(
1521                Some(OwnedRow::new(vec![Some(ScalarImpl::Int64(10))])),
1522                None,
1523                10,
1524                false,
1525            )
1526            .await
1527            .unwrap();
1528        state_writer
1529            .commit_state(Barrier::new_test_barrier(test_epoch(2)).epoch)
1530            .await
1531            .unwrap();
1532
1533        let mut restored_state = CdcBackfillState::new(
1534            TableId::new(1234),
1535            create_cdc_state_table(memory_state_store).await,
1536            5,
1537        );
1538        restored_state
1539            .init_epoch(Barrier::new_test_barrier(test_epoch(2)).epoch)
1540            .await
1541            .unwrap();
1542
1543        assert!(restored_state.restore_state().await.is_err());
1544    }
1545
1546    fn create_raw_cdc_chunk(rows: &[(&str, &str)]) -> StreamChunk {
1547        let schema = Schema::new(vec![
1548            Field::unnamed(DataType::Jsonb),
1549            Field::unnamed(DataType::Varchar),
1550        ]);
1551        let mut builders = schema.create_array_builders(rows.len());
1552        for (payload, offset) in rows {
1553            let payload_datum: Datum = Some(JsonbVal::from_str(payload).unwrap().into());
1554            let offset_datum: Datum = Some((*offset).into());
1555            builders[0].append(payload_datum);
1556            builders[1].append(offset_datum);
1557        }
1558        let columns = builders
1559            .into_iter()
1560            .map(|builder| builder.finish().into())
1561            .collect();
1562        StreamChunk::from_parts(
1563            vec![Op::Insert; rows.len()],
1564            DataChunk::new(columns, rows.len()),
1565        )
1566    }
1567
1568    async fn create_cdc_state_table(store: MemoryStateStore) -> StateTable<MemoryStateStore> {
1569        let state_schema = Schema::new(vec![
1570            Field::with_name(DataType::Varchar, "split_id"),
1571            Field::with_name(DataType::Int64, "id"),
1572            Field::with_name(DataType::Boolean, "backfill_finished"),
1573            Field::with_name(DataType::Int64, "row_count"),
1574            Field::with_name(DataType::Jsonb, "cdc_offset"),
1575        ]);
1576        let column_descs = vec![
1577            ColumnDesc::unnamed(ColumnId::from(0), state_schema[0].data_type.clone()),
1578            ColumnDesc::unnamed(ColumnId::from(1), state_schema[1].data_type.clone()),
1579            ColumnDesc::unnamed(ColumnId::from(2), state_schema[2].data_type.clone()),
1580            ColumnDesc::unnamed(ColumnId::from(3), state_schema[3].data_type.clone()),
1581            ColumnDesc::unnamed(ColumnId::from(4), state_schema[4].data_type.clone()),
1582        ];
1583
1584        StateTable::from_table_catalog(
1585            &gen_pbtable(
1586                TableId::from(0x42),
1587                column_descs,
1588                vec![OrderType::ascending()],
1589                vec![0],
1590                0,
1591            ),
1592            store,
1593            None,
1594        )
1595        .await
1596    }
1597
1598    async fn create_recovering_cdc_backfill() -> (
1599        MessageSender,
1600        CdcBackfillExecutor<MemoryStateStore>,
1601        MemoryStateStore,
1602    ) {
1603        let memory_state_store = MemoryStateStore::new();
1604        let mut state_writer = CdcBackfillState::new(
1605            TableId::new(1234),
1606            create_cdc_state_table(memory_state_store.clone()).await,
1607            5,
1608        );
1609        state_writer
1610            .init_epoch(Barrier::new_test_barrier(test_epoch(1)).epoch)
1611            .await
1612            .unwrap();
1613        state_writer
1614            .mutate_state(
1615                Some(OwnedRow::new(vec![Some(ScalarImpl::Int64(5))])),
1616                Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 2))),
1617                5,
1618                false,
1619            )
1620            .await
1621            .unwrap();
1622        state_writer
1623            .commit_state(Barrier::new_test_barrier(test_epoch(2)).epoch)
1624            .await
1625            .unwrap();
1626
1627        let (tx, source) = MockSource::channel();
1628        let source = source.into_executor(
1629            Schema::new(vec![
1630                Field::unnamed(DataType::Jsonb),
1631                Field::unnamed(DataType::Varchar),
1632            ]),
1633            vec![0],
1634        );
1635        let external_table = ExternalStorageTable::new(
1636            TableId::new(1234),
1637            SchemaTableName {
1638                schema_name: "public".to_owned(),
1639                table_name: "mock_table".to_owned(),
1640            },
1641            "mydb".to_owned(),
1642            ExternalTableConfig::default(),
1643            ExternalCdcTableType::Mock,
1644            Schema::new(vec![
1645                Field::with_name(DataType::Int64, "id"),
1646                Field::with_name(DataType::Float64, "price"),
1647            ]),
1648            vec![OrderType::ascending()],
1649            Some(vec![CdcKeyComparison::Native]),
1650            vec![0],
1651        );
1652        let output_columns = vec![
1653            ColumnDesc::named("id", ColumnId::new(1), DataType::Int64),
1654            ColumnDesc::named("price", ColumnId::new(2), DataType::Float64),
1655        ];
1656        let executor = CdcBackfillExecutor::new(
1657            ActorContext::for_test(0x1a),
1658            external_table,
1659            source,
1660            vec![0, 1],
1661            output_columns,
1662            None,
1663            StreamingMetrics::unused().into(),
1664            create_cdc_state_table(memory_state_store.clone()).await,
1665            None,
1666            CdcScanOptions::default(),
1667            BTreeMap::default(),
1668        );
1669        (tx, executor, memory_state_store)
1670    }
1671
1672    async fn assert_recovery_offset(
1673        memory_state_store: MemoryStateStore,
1674        expected_position: u64,
1675        epoch: u64,
1676    ) {
1677        let mut restored_state = CdcBackfillState::new(
1678            TableId::new(1234),
1679            create_cdc_state_table(memory_state_store).await,
1680            5,
1681        );
1682        restored_state
1683            .init_epoch(Barrier::new_test_barrier(epoch).epoch)
1684            .await
1685            .unwrap();
1686        let state = restored_state.restore_state().await.unwrap();
1687        assert_eq!(
1688            state.last_cdc_offset,
1689            Some(CdcOffset::MySql(MySqlOffset::new(
1690                "1.binlog".to_owned(),
1691                expected_position,
1692            )))
1693        );
1694    }
1695
1696    #[tokio::test]
1697    async fn test_recovery_filters_upstream_while_reader_is_pending() {
1698        let chunk = StreamChunk::from_rows(
1699            &[
1700                (
1701                    Op::Insert,
1702                    OwnedRow::new(vec![
1703                        Some(ScalarImpl::Int64(4)),
1704                        Some(ScalarImpl::Float64(44.04.into())),
1705                        Some(ScalarImpl::Utf8(
1706                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#
1707                                .into(),
1708                        )),
1709                    ]),
1710                ),
1711                (
1712                    Op::Insert,
1713                    OwnedRow::new(vec![
1714                        Some(ScalarImpl::Int64(6)),
1715                        Some(ScalarImpl::Float64(66.06.into())),
1716                        Some(ScalarImpl::Utf8(
1717                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
1718                                .into(),
1719                        )),
1720                    ]),
1721                ),
1722            ],
1723            &[DataType::Int64, DataType::Float64, DataType::Varchar],
1724        );
1725        let mut upstream = stream::iter([
1726            Ok(Message::Chunk(chunk)),
1727            Ok(Message::Barrier(Barrier::new_test_barrier(test_epoch(4)))),
1728        ]);
1729        let mut table_reader = None;
1730        let mut reader_future = Box::pin(std::future::pending::<ExternalTableReaderImpl>());
1731
1732        let Message::Chunk(chunk) =
1733            build_reader_and_poll_upstream(&mut upstream, &mut table_reader, &mut reader_future)
1734                .await
1735                .unwrap()
1736                .unwrap()
1737        else {
1738            panic!("expected an upstream chunk while the reader is pending");
1739        };
1740        assert!(table_reader.is_none());
1741
1742        let parser = ExternalCdcTableType::Mock.get_cdc_offset_parser().unwrap();
1743        let (chunk, consumed_offset) =
1744            CdcBackfillExecutor::<MemoryStateStore>::filter_recovery_chunk(
1745                &parser,
1746                chunk,
1747                &OwnedRow::new(vec![Some(ScalarImpl::Int64(5))]),
1748                PkCompareInfo {
1749                    indices: &[0],
1750                    order: &[OrderType::ascending()],
1751                    needs_unsigned_i64_compare: &[false],
1752                },
1753                &Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 2))),
1754                &[0, 1],
1755            )
1756            .unwrap();
1757        let chunk = chunk.expect("the scanned row should be emitted");
1758        assert_eq!(chunk.cardinality(), 1);
1759        assert_eq!(
1760            chunk.row_at(0).1.to_owned_row()[0],
1761            Some(ScalarImpl::Int64(4))
1762        );
1763        assert_eq!(
1764            consumed_offset,
1765            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 4)))
1766        );
1767        assert!(matches!(
1768            build_reader_and_poll_upstream(&mut upstream, &mut table_reader, &mut reader_future)
1769                .await
1770                .unwrap(),
1771            Some(Message::Barrier(_))
1772        ));
1773    }
1774
1775    #[tokio::test]
1776    async fn test_recovery_filters_updates_before_backfill_barrier() {
1777        let (mut tx, executor, memory_state_store) = create_recovering_cdc_backfill().await;
1778        let executor = executor.execute_inner();
1779        pin_mut!(executor);
1780
1781        tx.send_barrier(Barrier::new_test_barrier(test_epoch(3)));
1782        assert!(matches!(
1783            executor.next().await.unwrap().unwrap(),
1784            Message::Barrier(_)
1785        ));
1786
1787        tx.push_chunk(create_raw_cdc_chunk(&[
1788            (
1789                r#"{ "payload": { "before": { "id": 1, "price": 11.01 }, "after": { "id": 1, "price": 12.01 }, "source": { "version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002" }, "op": "u", "ts_ms": 1695277757017, "transaction": null } }"#,
1790                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#,
1791            ),
1792            (
1793                r#"{ "payload": { "before": { "id": 6, "price": 66.06 }, "after": { "id": 6, "price": 67.06 }, "source": { "version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002" }, "op": "u", "ts_ms": 1695277757017, "transaction": null } }"#,
1794                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#,
1795            ),
1796        ]));
1797        tx.send_barrier(Barrier::new_test_barrier(test_epoch(4)));
1798
1799        let Message::Chunk(chunk) = executor.next().await.unwrap().unwrap() else {
1800            panic!("expected the recovered update before the barrier");
1801        };
1802        let rows = chunk
1803            .rows()
1804            .map(|(op, row)| (op, row.to_owned_row()))
1805            .collect::<Vec<_>>();
1806        assert_eq!(rows.len(), 1);
1807        assert_eq!(rows[0].0, Op::Insert);
1808        assert_eq!(rows[0].1[0], Some(ScalarImpl::Int64(1)));
1809        assert!(matches!(
1810            executor.next().await.unwrap().unwrap(),
1811            Message::Barrier(_)
1812        ));
1813        assert_recovery_offset(memory_state_store, 2, test_epoch(4)).await;
1814    }
1815
1816    #[tokio::test]
1817    async fn test_recovery_preserves_scanned_inserts_and_deletes() {
1818        let (mut tx, executor, _) = create_recovering_cdc_backfill().await;
1819        let executor = executor.execute_inner();
1820        pin_mut!(executor);
1821
1822        tx.send_barrier(Barrier::new_test_barrier(test_epoch(3)));
1823        assert!(matches!(
1824            executor.next().await.unwrap().unwrap(),
1825            Message::Barrier(_)
1826        ));
1827        // This replayed event predates the saved offset and must not lower the offset filter.
1828        tx.push_chunk(create_raw_cdc_chunk(&[(
1829            r#"{ "payload": { "before": null, "after": { "id": 3, "price": 30.01 }, "op": "c" } }"#,
1830            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":1},"isHeartbeat":false}"#,
1831        )]));
1832        tx.push_chunk(create_raw_cdc_chunk(&[
1833            (
1834                r#"{ "payload": { "before": null, "after": { "id": 3, "price": 31.01 }, "op": "c" } }"#,
1835                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":1},"isHeartbeat":false}"#,
1836            ),
1837            (
1838                r#"{ "payload": { "before": null, "after": { "id": 4, "price": 44.04 }, "op": "c" } }"#,
1839                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":2},"isHeartbeat":false}"#,
1840            ),
1841            (
1842                r#"{ "payload": { "before": { "id": 5, "price": 55.05 }, "after": null, "op": "d" } }"#,
1843                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#,
1844            ),
1845        ]));
1846        tx.send_barrier(Barrier::new_test_barrier(test_epoch(4)));
1847
1848        let Message::Chunk(chunk) = executor.next().await.unwrap().unwrap() else {
1849            panic!("expected the scanned-prefix changes before the barrier");
1850        };
1851        let rows = chunk
1852            .rows()
1853            .map(|(op, row)| (op, row.to_owned_row()))
1854            .collect::<Vec<_>>();
1855        assert_eq!(rows.len(), 2);
1856        assert_eq!(rows[0].0, Op::Insert);
1857        assert_eq!(rows[0].1[0], Some(ScalarImpl::Int64(4)));
1858        assert_eq!(rows[1].0, Op::Delete);
1859        assert_eq!(rows[1].1[0], Some(ScalarImpl::Int64(5)));
1860        assert!(matches!(
1861            executor.next().await.unwrap().unwrap(),
1862            Message::Barrier(_)
1863        ));
1864    }
1865
1866    #[tokio::test]
1867    async fn test_recovery_does_not_leave_deleted_unscanned_row() {
1868        let (mut tx, mut executor, memory_state_store) = create_recovering_cdc_backfill().await;
1869        executor.options.snapshot_barrier_interval = 1;
1870        let executor = executor.execute_inner();
1871        pin_mut!(executor);
1872
1873        tx.send_barrier(Barrier::new_test_barrier(test_epoch(3)));
1874        assert!(matches!(
1875            executor.next().await.unwrap().unwrap(),
1876            Message::Barrier(_)
1877        ));
1878
1879        // Recover at PK 5, then receive an insert for an unscanned key during startup.
1880        tx.push_chunk(create_raw_cdc_chunk(&[(
1881            r#"{ "payload": { "before": null, "after": { "id": 100, "price": 100.01 }, "op": "c" } }"#,
1882            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#,
1883        )]));
1884        tx.send_barrier(Barrier::new_test_barrier(test_epoch(4)));
1885
1886        // Normal backfill receives its delete before the snapshot reaches PK 100.
1887        // The snapshot refresh at barrier 5 scans up to PK 8 and filters out the delete.
1888        tx.push_chunk(create_raw_cdc_chunk(&[(
1889            r#"{ "payload": { "before": { "id": 100, "price": 100.01 }, "after": null, "op": "d" } }"#,
1890            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#,
1891        )]));
1892        tx.send_barrier(Barrier::new_test_barrier(test_epoch(5)));
1893        // The next snapshot is empty, so backfill completes at barrier 6.
1894        tx.send_barrier(Barrier::new_test_barrier(test_epoch(6)));
1895
1896        let mut materialized = BTreeMap::new();
1897        loop {
1898            match executor.next().await.unwrap().unwrap() {
1899                Message::Chunk(chunk) => {
1900                    for (op, row) in chunk.rows() {
1901                        let pk = row.datum_at(0).unwrap().into_int64();
1902                        match op {
1903                            Op::Insert => {
1904                                materialized.insert(pk, row.to_owned_row());
1905                            }
1906                            Op::Delete => {
1907                                materialized.remove(&pk);
1908                            }
1909                            _ => unreachable!("CDC updates are converted to inserts"),
1910                        }
1911                    }
1912                }
1913                Message::Barrier(barrier) if barrier.epoch.curr == test_epoch(6) => break,
1914                Message::Barrier(_) => {}
1915                Message::Watermark(_) => panic!("unexpected watermark during backfill"),
1916            }
1917        }
1918
1919        assert!(!materialized.contains_key(&100));
1920        let mut restored_state = CdcBackfillState::new(
1921            TableId::new(1234),
1922            create_cdc_state_table(memory_state_store).await,
1923            5,
1924        );
1925        restored_state
1926            .init_epoch(Barrier::new_test_barrier(test_epoch(6)).epoch)
1927            .await
1928            .unwrap();
1929        assert!(restored_state.restore_state().await.unwrap().is_finished);
1930    }
1931
1932    #[test]
1933    fn test_consume_upstream_chunk_buffer_retains_future_rows() {
1934        let mut upstream_chunk_buffer = vec![StreamChunk::from_rows(
1935            &[
1936                (
1937                    Op::Insert,
1938                    OwnedRow::new(vec![
1939                        Some(ScalarImpl::Int64(1)),
1940                        Some(ScalarImpl::Int64(100)),
1941                        Some(ScalarImpl::Utf8(
1942                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#
1943                                .into(),
1944                        )),
1945                    ]),
1946                ),
1947                (
1948                    Op::Insert,
1949                    OwnedRow::new(vec![
1950                        Some(ScalarImpl::Int64(6)),
1951                        Some(ScalarImpl::Int64(600)),
1952                        Some(ScalarImpl::Utf8(
1953                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
1954                                .into(),
1955                        )),
1956                    ]),
1957                ),
1958            ],
1959            &[DataType::Int64, DataType::Int64, DataType::Varchar],
1960        )];
1961
1962        let (emitted_chunks, drained_row_count, drained_offset) =
1963            CdcBackfillExecutor::<MemoryStateStore>::consume_upstream_chunk_buffer(
1964                &MockExternalTableReader::get_cdc_offset_parser(),
1965                &mut upstream_chunk_buffer,
1966                Some(&OwnedRow::new(vec![Some(ScalarImpl::Int64(5))])),
1967                PkCompareInfo {
1968                    indices: &[0],
1969                    order: &[OrderType::ascending()],
1970                    needs_unsigned_i64_compare: &[false],
1971                },
1972                &Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 2))),
1973                &[0, 1],
1974            )
1975            .unwrap();
1976
1977        assert_eq!(drained_row_count, 1);
1978        assert_eq!(
1979            drained_offset,
1980            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 3)))
1981        );
1982        assert_eq!(emitted_chunks.len(), 1);
1983        assert_eq!(emitted_chunks[0].rows().count(), 1);
1984        assert_eq!(
1985            emitted_chunks[0].rows().next().unwrap().1.to_owned_row(),
1986            OwnedRow::new(vec![
1987                Some(ScalarImpl::Int64(1)),
1988                Some(ScalarImpl::Int64(100))
1989            ])
1990        );
1991
1992        assert_eq!(upstream_chunk_buffer.len(), 1);
1993        assert_eq!(upstream_chunk_buffer[0].rows().count(), 1);
1994        assert_eq!(
1995            upstream_chunk_buffer[0]
1996                .rows()
1997                .next()
1998                .unwrap()
1999                .1
2000                .to_owned_row(),
2001            OwnedRow::new(vec![
2002                Some(ScalarImpl::Int64(6)),
2003                Some(ScalarImpl::Int64(600)),
2004                Some(ScalarImpl::Utf8(
2005                    r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2006                        .into(),
2007                )),
2008            ])
2009        );
2010    }
2011
2012    #[test]
2013    fn test_consume_buffer_uses_unsigned_bigint_pk_order() {
2014        let mut upstream_chunk_buffer = vec![StreamChunk::from_rows(
2015            &[
2016                (
2017                    Op::Insert,
2018                    OwnedRow::new(vec![
2019                        Some(ScalarImpl::Int64(4)),
2020                        Some(ScalarImpl::Int64(400)),
2021                        Some(ScalarImpl::Utf8(
2022                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#
2023                                .into(),
2024                        )),
2025                    ]),
2026                ),
2027                (
2028                    Op::Insert,
2029                    OwnedRow::new(vec![
2030                        // `u64::MAX` represented in RisingWave's `i64` storage.
2031                        Some(ScalarImpl::Int64(-1)),
2032                        Some(ScalarImpl::Int64(900)),
2033                        Some(ScalarImpl::Utf8(
2034                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2035                                .into(),
2036                        )),
2037                    ]),
2038                ),
2039            ],
2040            &[DataType::Int64, DataType::Int64, DataType::Varchar],
2041        )];
2042
2043        let (emitted_chunks, drained_row_count, drained_offset) =
2044            CdcBackfillExecutor::<MemoryStateStore>::consume_upstream_chunk_buffer(
2045                &MockExternalTableReader::get_cdc_offset_parser(),
2046                &mut upstream_chunk_buffer,
2047                Some(&OwnedRow::new(vec![Some(ScalarImpl::Int64(5))])),
2048                PkCompareInfo {
2049                    indices: &[0],
2050                    order: &[OrderType::ascending()],
2051                    needs_unsigned_i64_compare: &[true],
2052                },
2053                &Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 2))),
2054                &[0, 1],
2055            )
2056            .unwrap();
2057
2058        assert_eq!(drained_row_count, 1);
2059        assert_eq!(
2060            drained_offset,
2061            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 3)))
2062        );
2063        assert_eq!(emitted_chunks.len(), 1);
2064        assert_eq!(
2065            emitted_chunks[0].rows().next().unwrap().1.to_owned_row(),
2066            OwnedRow::new(vec![
2067                Some(ScalarImpl::Int64(4)),
2068                Some(ScalarImpl::Int64(400))
2069            ])
2070        );
2071
2072        assert_eq!(upstream_chunk_buffer.len(), 1);
2073        assert_eq!(
2074            upstream_chunk_buffer[0]
2075                .rows()
2076                .next()
2077                .unwrap()
2078                .1
2079                .to_owned_row(),
2080            OwnedRow::new(vec![
2081                Some(ScalarImpl::Int64(-1)),
2082                Some(ScalarImpl::Int64(900)),
2083                Some(ScalarImpl::Utf8(
2084                    r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2085                        .into(),
2086                )),
2087            ])
2088        );
2089    }
2090
2091    #[test]
2092    fn test_consume_buffer_non_monotonic_pk_in_chunk() {
2093        let mut upstream_chunk_buffer = vec![StreamChunk::from_rows(
2094            &[
2095                (
2096                    Op::Insert,
2097                    OwnedRow::new(vec![
2098                        Some(ScalarImpl::Int64(1)),
2099                        Some(ScalarImpl::Int64(100)),
2100                        Some(ScalarImpl::Utf8(
2101                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#
2102                                .into(),
2103                        )),
2104                    ]),
2105                ),
2106                (
2107                    Op::Insert,
2108                    OwnedRow::new(vec![
2109                        Some(ScalarImpl::Int64(6)),
2110                        Some(ScalarImpl::Int64(600)),
2111                        Some(ScalarImpl::Utf8(
2112                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2113                                .into(),
2114                        )),
2115                    ]),
2116                ),
2117                (
2118                    Op::Insert,
2119                    OwnedRow::new(vec![
2120                        Some(ScalarImpl::Int64(2)),
2121                        Some(ScalarImpl::Int64(200)),
2122                        Some(ScalarImpl::Utf8(
2123                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":5},"isHeartbeat":false}"#
2124                                .into(),
2125                        )),
2126                    ]),
2127                ),
2128            ],
2129            &[DataType::Int64, DataType::Int64, DataType::Varchar],
2130        )];
2131
2132        let (emitted_chunks, drained_row_count, drained_offset) =
2133            CdcBackfillExecutor::<MemoryStateStore>::consume_upstream_chunk_buffer(
2134                &MockExternalTableReader::get_cdc_offset_parser(),
2135                &mut upstream_chunk_buffer,
2136                Some(&OwnedRow::new(vec![Some(ScalarImpl::Int64(5))])),
2137                PkCompareInfo {
2138                    indices: &[0],
2139                    order: &[OrderType::ascending()],
2140                    needs_unsigned_i64_compare: &[false],
2141                },
2142                &Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 2))),
2143                &[0, 1],
2144            )
2145            .unwrap();
2146
2147        assert_eq!(drained_row_count, 2);
2148        assert_eq!(
2149            drained_offset,
2150            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 3)))
2151        );
2152        assert_eq!(emitted_chunks.len(), 1);
2153        assert_eq!(emitted_chunks[0].rows().count(), 2);
2154        assert_eq!(
2155            emitted_chunks[0]
2156                .rows()
2157                .map(|(_, row)| row.to_owned_row())
2158                .collect::<Vec<_>>(),
2159            vec![
2160                OwnedRow::new(vec![
2161                    Some(ScalarImpl::Int64(1)),
2162                    Some(ScalarImpl::Int64(100))
2163                ]),
2164                OwnedRow::new(vec![
2165                    Some(ScalarImpl::Int64(2)),
2166                    Some(ScalarImpl::Int64(200))
2167                ]),
2168            ]
2169        );
2170
2171        assert_eq!(upstream_chunk_buffer.len(), 1);
2172        assert_eq!(upstream_chunk_buffer[0].rows().count(), 1);
2173        assert_eq!(
2174            upstream_chunk_buffer[0]
2175                .rows()
2176                .next()
2177                .unwrap()
2178                .1
2179                .to_owned_row(),
2180            OwnedRow::new(vec![
2181                Some(ScalarImpl::Int64(6)),
2182                Some(ScalarImpl::Int64(600)),
2183                Some(ScalarImpl::Utf8(
2184                    r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2185                        .into(),
2186                )),
2187            ])
2188        );
2189    }
2190
2191    #[test]
2192    fn test_consume_buffer_processes_chunks_after_future_row() {
2193        let mut upstream_chunk_buffer = vec![
2194            StreamChunk::from_rows(
2195                &[
2196                    (
2197                        Op::Insert,
2198                        OwnedRow::new(vec![
2199                            Some(ScalarImpl::Int64(1)),
2200                            Some(ScalarImpl::Int64(100)),
2201                            Some(ScalarImpl::Utf8(
2202                                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#
2203                                    .into(),
2204                            )),
2205                        ]),
2206                    ),
2207                    (
2208                        Op::Insert,
2209                        OwnedRow::new(vec![
2210                            Some(ScalarImpl::Int64(6)),
2211                            Some(ScalarImpl::Int64(600)),
2212                            Some(ScalarImpl::Utf8(
2213                                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2214                                    .into(),
2215                            )),
2216                        ]),
2217                    ),
2218                ],
2219                &[DataType::Int64, DataType::Int64, DataType::Varchar],
2220            ),
2221            StreamChunk::from_rows(
2222                &[(
2223                    Op::Insert,
2224                    OwnedRow::new(vec![
2225                        Some(ScalarImpl::Int64(2)),
2226                        Some(ScalarImpl::Int64(200)),
2227                        Some(ScalarImpl::Utf8(
2228                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":5},"isHeartbeat":false}"#
2229                                .into(),
2230                        )),
2231                    ]),
2232                )],
2233                &[DataType::Int64, DataType::Int64, DataType::Varchar],
2234            ),
2235        ];
2236
2237        let (emitted_chunks, drained_row_count, drained_offset) =
2238            CdcBackfillExecutor::<MemoryStateStore>::consume_upstream_chunk_buffer(
2239                &MockExternalTableReader::get_cdc_offset_parser(),
2240                &mut upstream_chunk_buffer,
2241                Some(&OwnedRow::new(vec![Some(ScalarImpl::Int64(5))])),
2242                PkCompareInfo {
2243                    indices: &[0],
2244                    order: &[OrderType::ascending()],
2245                    needs_unsigned_i64_compare: &[false],
2246                },
2247                &Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 2))),
2248                &[0, 1],
2249            )
2250            .unwrap();
2251
2252        assert_eq!(drained_row_count, 2);
2253        assert_eq!(
2254            drained_offset,
2255            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 3)))
2256        );
2257        assert_eq!(emitted_chunks.len(), 2);
2258        assert_eq!(emitted_chunks[0].rows().count(), 1);
2259        assert_eq!(emitted_chunks[1].rows().count(), 1);
2260        assert_eq!(
2261            emitted_chunks[1].rows().next().unwrap().1.to_owned_row(),
2262            OwnedRow::new(vec![
2263                Some(ScalarImpl::Int64(2)),
2264                Some(ScalarImpl::Int64(200))
2265            ])
2266        );
2267
2268        assert_eq!(upstream_chunk_buffer.len(), 1);
2269        assert_eq!(
2270            upstream_chunk_buffer[0]
2271                .rows()
2272                .next()
2273                .unwrap()
2274                .1
2275                .to_owned_row(),
2276            OwnedRow::new(vec![
2277                Some(ScalarImpl::Int64(6)),
2278                Some(ScalarImpl::Int64(600)),
2279                Some(ScalarImpl::Utf8(
2280                    r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2281                        .into(),
2282                )),
2283            ])
2284        );
2285    }
2286
2287    #[test]
2288    fn test_consume_buffer_advances_offset_after_preceding_rows_are_emitted() {
2289        let mut upstream_chunk_buffer = vec![StreamChunk::from_rows(
2290            &[
2291                (
2292                    Op::Insert,
2293                    OwnedRow::new(vec![
2294                        Some(ScalarImpl::Int64(6)),
2295                        Some(ScalarImpl::Int64(600)),
2296                        Some(ScalarImpl::Utf8(
2297                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#
2298                                .into(),
2299                        )),
2300                    ]),
2301                ),
2302                (
2303                    Op::Insert,
2304                    OwnedRow::new(vec![
2305                        Some(ScalarImpl::Int64(2)),
2306                        Some(ScalarImpl::Int64(200)),
2307                        Some(ScalarImpl::Utf8(
2308                            r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":5},"isHeartbeat":false}"#
2309                                .into(),
2310                        )),
2311                    ]),
2312                ),
2313            ],
2314            &[DataType::Int64, DataType::Int64, DataType::Varchar],
2315        )];
2316
2317        let (_, drained_row_count, drained_offset) =
2318            CdcBackfillExecutor::<MemoryStateStore>::consume_upstream_chunk_buffer(
2319                &MockExternalTableReader::get_cdc_offset_parser(),
2320                &mut upstream_chunk_buffer,
2321                Some(&OwnedRow::new(vec![Some(ScalarImpl::Int64(6))])),
2322                PkCompareInfo {
2323                    indices: &[0],
2324                    order: &[OrderType::ascending()],
2325                    needs_unsigned_i64_compare: &[false],
2326                },
2327                &Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 3))),
2328                &[0, 1],
2329            )
2330            .unwrap();
2331
2332        assert_eq!(drained_row_count, 2);
2333        assert_eq!(
2334            drained_offset,
2335            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 5)))
2336        );
2337        assert!(upstream_chunk_buffer.is_empty());
2338    }
2339
2340    #[tokio::test]
2341    async fn test_cdc_backfill_persists_buffered_offset_on_checkpoint() {
2342        let memory_state_store = MemoryStateStore::new();
2343        let state_table = create_cdc_state_table(memory_state_store.clone()).await;
2344
2345        let (mut tx, source) = MockSource::channel();
2346        let source = source.into_executor(
2347            Schema::new(vec![
2348                Field::unnamed(DataType::Jsonb),
2349                Field::unnamed(DataType::Varchar),
2350            ]),
2351            vec![0],
2352        );
2353
2354        let external_table = ExternalStorageTable::new(
2355            TableId::new(1234),
2356            SchemaTableName {
2357                schema_name: "public".to_owned(),
2358                table_name: "mock_table".to_owned(),
2359            },
2360            "mydb".to_owned(),
2361            ExternalTableConfig::default(),
2362            ExternalCdcTableType::Mock,
2363            Schema::new(vec![
2364                Field::with_name(DataType::Int64, "id"),
2365                Field::with_name(DataType::Float64, "price"),
2366            ]),
2367            vec![OrderType::ascending()],
2368            Some(vec![CdcKeyComparison::Native]),
2369            vec![0],
2370        );
2371        let output_columns = vec![
2372            ColumnDesc::named("id", ColumnId::new(1), DataType::Int64),
2373            ColumnDesc::named("price", ColumnId::new(2), DataType::Float64),
2374        ];
2375
2376        let executor = CdcBackfillExecutor::new(
2377            ActorContext::for_test(0x1a),
2378            external_table,
2379            source,
2380            vec![0, 1],
2381            output_columns,
2382            None,
2383            StreamingMetrics::unused().into(),
2384            state_table,
2385            None,
2386            CdcScanOptions {
2387                snapshot_barrier_interval: 10,
2388                ..Default::default()
2389            },
2390            BTreeMap::default(),
2391        )
2392        .execute_inner();
2393        pin_mut!(executor);
2394
2395        tx.send_barrier(Barrier::new_test_barrier(test_epoch(1)));
2396        assert!(matches!(
2397            executor.next().await.unwrap().unwrap(),
2398            Message::Barrier(_)
2399        ));
2400
2401        tx.send_barrier(Barrier::new_test_barrier(test_epoch(2)));
2402        assert!(matches!(
2403            executor.next().await.unwrap().unwrap(),
2404            Message::Barrier(_)
2405        ));
2406
2407        assert!(matches!(
2408            executor.next().await.unwrap().unwrap(),
2409            Message::Chunk(_)
2410        ));
2411
2412        tx.push_chunk(create_raw_cdc_chunk(&[
2413            (
2414                r#"{ "payload": { "before": null, "after": { "id": 1, "price": 10.01 }, "source": { "version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002" }, "op": "r", "ts_ms": 1695277757017, "transaction": null } }"#,
2415                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":3},"isHeartbeat":false}"#,
2416            ),
2417            (
2418                r#"{ "payload": { "before": null, "after": { "id": 6, "price": 66.06 }, "source": { "version": "1.9.7.Final", "connector": "mysql", "name": "RW_CDC_1002" }, "op": "r", "ts_ms": 1695277757017, "transaction": null } }"#,
2419                r#"{"sourcePartition":{},"sourceOffset":{"file":"1.binlog","pos":4},"isHeartbeat":false}"#,
2420            ),
2421        ]));
2422        tx.send_barrier(Barrier::new_test_barrier(test_epoch(3)));
2423
2424        assert!(matches!(
2425            executor.next().await.unwrap().unwrap(),
2426            Message::Chunk(_)
2427        ));
2428        assert!(matches!(
2429            executor.next().await.unwrap().unwrap(),
2430            Message::Barrier(_)
2431        ));
2432
2433        let mut restored_state = CdcBackfillState::new(
2434            TableId::new(1234),
2435            create_cdc_state_table(memory_state_store).await,
2436            5,
2437        );
2438        restored_state
2439            .init_epoch(Barrier::new_test_barrier(test_epoch(3)).epoch)
2440            .await
2441            .unwrap();
2442        let state = restored_state.restore_state().await.unwrap();
2443        assert_eq!(
2444            state.last_cdc_offset,
2445            Some(CdcOffset::MySql(MySqlOffset::new("1.binlog".to_owned(), 4)))
2446        );
2447    }
2448}