1use 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
58const 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
74pub(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 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 external_table: ExternalStorageTable,
124
125 upstream: Executor,
127
128 output_indices: Vec<usize>,
130
131 output_columns: Vec<ColumnDesc>,
133
134 state_impl: CdcBackfillState<S>,
136
137 progress: Option<CreateMviewProgressReporter>,
140
141 metrics: CdcBackfillMetrics,
142
143 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 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 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 let mut current_pk_pos: Option<OwnedRow>;
358
359 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 yield Message::Barrier(first_barrier);
366
367 let mut state_impl = self.state_impl;
370
371 state_impl.init_epoch(first_barrier_epoch).await?;
372
373 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 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 let handle_toast_columns: bool =
386 self.external_table.table_type() == &ExternalCdcTableType::Postgres;
387 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 if need_backfill {
422 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 }
494 }
495 Message::Watermark(_) => {
496 }
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 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 let _ = Pin::new(&mut upstream).peek().await;
563
564 #[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 }
589 }
590 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 }
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 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 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 Either::Left(msg) => {
681 match msg? {
682 Message::Barrier(barrier) => {
683 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 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 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 if can_start_new_snapshot || needs_rebuild_snapshot {
733 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;
743 } else {
744 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 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 yield Message::Barrier(barrier);
792 }
793 }
794 Message::Chunk(chunk) => {
795 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 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 upstream_chunk_buffer.push(chunk.compact_vis());
825 }
826 Message::Watermark(_) => {
827 }
829 }
830 }
831 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 for chunk in upstream_chunk_buffer.drain(..) {
846 yield Message::Chunk(mapping_chunk(
847 chunk,
848 &self.output_indices,
849 ));
850 }
851
852 break 'backfill_loop;
854 }
855 Some(chunk) => {
856 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 let (_, mut snapshot_stream) = backfill_stream.into_inner();
888 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 for chunk in upstream_chunk_buffer.drain(..) {
912 yield Message::Chunk(mapping_chunk(chunk, &self.output_indices));
913 }
914
915 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 state_impl.commit_state(pending_barrier.epoch).await?;
927 yield Message::Barrier(pending_barrier);
928 break 'backfill_loop;
930 }
931 Some(_) if is_snapshot_paused => {
932 }
934 Some(chunk) => {
935 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 let Some(current_pos) = ¤t_pk_pos {
957 for chunk in upstream_chunk_buffer.drain(..) {
958 cur_barrier_upstream_processed_rows += chunk.cardinality() as u64;
959
960 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 upstream_chunk_buffer.clear();
982 }
983
984 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 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 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 while let Some(Ok(msg)) = upstream.next().await {
1035 if let Some(msg) = mapping_message(msg, &self.output_indices) {
1036 if let Message::Barrier(barrier) = &msg {
1038 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 if let Some(progress) = self.progress.as_mut() {
1052 progress.finish(barrier.epoch, total_snapshot_row_count);
1053 }
1054 yield msg;
1055 break;
1057 }
1058 yield msg;
1059 }
1060 }
1061
1062 #[for_await]
1066 for msg in upstream {
1067 if let Some(msg) = mapping_message(msg?, &self.output_indices) {
1070 if let Message::Barrier(barrier) = &msg {
1071 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 protocol_config: ProtocolProperties::Debezium(DebeziumProps::default()),
1121 };
1122
1123 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 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 let payloads = chunk.data_chunk().project(&[0]);
1170 let offsets = chunk.data_chunk().project(&[1]).compact_vis();
1171
1172 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()); let (ops, mut columns, vis) = parsed_chunk.into_inner();
1197 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), Field::unnamed(DataType::Varchar), Field::unnamed(DataType::Varchar), ]);
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#"{ "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 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 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), Field::unnamed(DataType::Varchar), Field::unnamed(DataType::Varchar), ]);
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]; 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 disable_backfill: true,
1365 ..CdcScanOptions::default()
1366 },
1367 BTreeMap::default(),
1368 );
1369 let s = cdc.execute_inner();
1373 pin_mut!(s);
1374
1375 tx.send_barrier(Barrier::new_test_barrier(test_epoch(8)));
1377 {
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 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 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 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 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 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 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}