1use std::collections::VecDeque;
62use std::future::pending;
63use std::mem::replace;
64use std::pin::Pin;
65
66use anyhow::anyhow;
67use futures::future::{BoxFuture, Either, select};
68use futures::stream::StreamFuture;
69use futures::{FutureExt, StreamExt, TryStreamExt};
70use futures_async_stream::try_stream;
71use risingwave_common::array::StreamChunk;
72use risingwave_common::bitmap::Bitmap;
73use risingwave_common::catalog::{TableId, TableOption};
74use risingwave_common::must_match;
75use risingwave_common::util::epoch::EpochPair;
76use risingwave_common_estimate_size::EstimateSize;
77use risingwave_connector::sink::log_store::{ChunkId, LogStoreResult};
78use risingwave_pb::id::FragmentId;
79use risingwave_storage::StateStore;
80use risingwave_storage::store::timeout_auto_rebuild::TimeoutAutoRebuildIter;
81use risingwave_storage::store::{
82 LocalStateStore, NewLocalOptions, OpConsistencyLevel, StateStoreRead,
83};
84use rw_futures_util::drop_either_future;
85use tokio::time::{Duration, Instant, Sleep, sleep_until};
86use tokio_stream::adapters::Peekable;
87
88use crate::common::log_store_impl::kv_log_store::buffer::LogStoreBufferItem;
89use crate::common::log_store_impl::kv_log_store::reader::LogStoreReadStateStreamRangeStart;
90use crate::common::log_store_impl::kv_log_store::serde::{
91 KvLogStoreItem, LogStoreItemMergeStream, LogStoreRowSerde,
92};
93use crate::common::log_store_impl::kv_log_store::state::{
94 LogStorePostSealCurrentEpoch, LogStoreReadState, LogStoreStateWriteChunkFuture,
95 LogStoreWriteState, new_log_store_state,
96};
97use crate::common::log_store_impl::kv_log_store::{
98 Epoch, FIRST_SEQ_ID, FlushInfo, LogStoreVnodeProgress, SeqId,
99};
100use crate::executor::prelude::*;
101use crate::executor::sync_kv_log_store::metrics::SyncedKvLogStoreMetrics;
102use crate::executor::{
103 Barrier, BoxedMessageStream, Message, StreamExecutorError, StreamExecutorResult, Watermark,
104};
105
106pub mod metrics {
107 use risingwave_common::id::FragmentId;
108 use risingwave_common::metrics::{LabelGuardedIntCounter, LabelGuardedIntGauge};
109
110 use crate::common::log_store_impl::kv_log_store::KvLogStoreReadMetrics;
111 use crate::executor::monitor::StreamingMetrics;
112 use crate::task::ActorId;
113
114 #[derive(Clone)]
115 pub struct SyncedKvLogStoreMetrics {
116 pub unclean_state: LabelGuardedIntCounter,
118 pub clean_state: LabelGuardedIntCounter,
119 pub wait_next_poll_ns: LabelGuardedIntCounter,
120
121 pub storage_write_count: LabelGuardedIntCounter,
123 pub storage_write_size: LabelGuardedIntCounter,
124 pub pause_duration_ns: LabelGuardedIntCounter,
125
126 pub buffer_unconsumed_item_count: LabelGuardedIntGauge,
128 pub buffer_unconsumed_row_count: LabelGuardedIntGauge,
129 pub buffer_unconsumed_epoch_count: LabelGuardedIntGauge,
130 pub buffer_unconsumed_min_epoch: LabelGuardedIntGauge,
131 pub buffer_memory_bytes: LabelGuardedIntGauge,
132 pub buffer_read_count: LabelGuardedIntCounter,
133 pub buffer_read_size: LabelGuardedIntCounter,
134
135 pub total_read_count: LabelGuardedIntCounter,
137 pub total_read_size: LabelGuardedIntCounter,
138 pub persistent_log_read_metrics: KvLogStoreReadMetrics,
139 pub flushed_buffer_read_metrics: KvLogStoreReadMetrics,
140 }
141
142 impl SyncedKvLogStoreMetrics {
143 pub(crate) fn new(
150 metrics: &StreamingMetrics,
151 actor_id: ActorId,
152 fragment_id: FragmentId,
153 name: &str,
154 target: &'static str,
155 ) -> Self {
156 let actor_id_str = actor_id.to_string();
157 let fragment_id_str = fragment_id.to_string();
158 let labels = &[&actor_id_str, target, &fragment_id_str, name];
159
160 let unclean_state = metrics.sync_kv_log_store_state.with_guarded_label_values(&[
161 "dirty",
162 &actor_id_str,
163 target,
164 &fragment_id_str,
165 name,
166 ]);
167 let clean_state = metrics.sync_kv_log_store_state.with_guarded_label_values(&[
168 "clean",
169 &actor_id_str,
170 target,
171 &fragment_id_str,
172 name,
173 ]);
174 let wait_next_poll_ns = metrics
175 .sync_kv_log_store_wait_next_poll_ns
176 .with_guarded_label_values(labels);
177
178 let storage_write_size = metrics
179 .sync_kv_log_store_storage_write_size
180 .with_guarded_label_values(labels);
181 let storage_write_count = metrics
182 .sync_kv_log_store_storage_write_count
183 .with_guarded_label_values(labels);
184 let storage_pause_duration_ns = metrics
185 .sync_kv_log_store_write_pause_duration_ns
186 .with_guarded_label_values(labels);
187
188 let buffer_unconsumed_item_count = metrics
189 .sync_kv_log_store_buffer_unconsumed_item_count
190 .with_guarded_label_values(labels);
191 let buffer_unconsumed_row_count = metrics
192 .sync_kv_log_store_buffer_unconsumed_row_count
193 .with_guarded_label_values(labels);
194 let buffer_unconsumed_epoch_count = metrics
195 .sync_kv_log_store_buffer_unconsumed_epoch_count
196 .with_guarded_label_values(labels);
197 let buffer_unconsumed_min_epoch = metrics
198 .sync_kv_log_store_buffer_unconsumed_min_epoch
199 .with_guarded_label_values(labels);
200 let buffer_memory_bytes = metrics
201 .sync_kv_log_store_buffer_memory_bytes
202 .with_guarded_label_values(labels);
203 let buffer_read_count = metrics
204 .sync_kv_log_store_read_count
205 .with_guarded_label_values(&[
206 "buffer",
207 &actor_id_str,
208 target,
209 &fragment_id_str,
210 name,
211 ]);
212
213 let buffer_read_size = metrics
214 .sync_kv_log_store_read_size
215 .with_guarded_label_values(&[
216 "buffer",
217 &actor_id_str,
218 target,
219 &fragment_id_str,
220 name,
221 ]);
222
223 let total_read_count = metrics
224 .sync_kv_log_store_read_count
225 .with_guarded_label_values(&[
226 "total",
227 &actor_id_str,
228 target,
229 &fragment_id_str,
230 name,
231 ]);
232
233 let total_read_size = metrics
234 .sync_kv_log_store_read_size
235 .with_guarded_label_values(&[
236 "total",
237 &actor_id_str,
238 target,
239 &fragment_id_str,
240 name,
241 ]);
242
243 const READ_PERSISTENT_LOG: &str = "persistent_log";
244 const READ_FLUSHED_BUFFER: &str = "flushed_buffer";
245
246 let persistent_log_read_size = metrics
247 .sync_kv_log_store_read_size
248 .with_guarded_label_values(&[
249 READ_PERSISTENT_LOG,
250 &actor_id_str,
251 target,
252 &fragment_id_str,
253 name,
254 ]);
255
256 let persistent_log_read_count = metrics
257 .sync_kv_log_store_read_count
258 .with_guarded_label_values(&[
259 READ_PERSISTENT_LOG,
260 &actor_id_str,
261 target,
262 &fragment_id_str,
263 name,
264 ]);
265
266 let flushed_buffer_read_size = metrics
267 .sync_kv_log_store_read_size
268 .with_guarded_label_values(&[
269 READ_FLUSHED_BUFFER,
270 &actor_id_str,
271 target,
272 &fragment_id_str,
273 name,
274 ]);
275
276 let flushed_buffer_read_count = metrics
277 .sync_kv_log_store_read_count
278 .with_guarded_label_values(&[
279 READ_FLUSHED_BUFFER,
280 &actor_id_str,
281 target,
282 &fragment_id_str,
283 name,
284 ]);
285
286 Self {
287 unclean_state,
288 clean_state,
289 wait_next_poll_ns,
290 storage_write_size,
291 storage_write_count,
292 pause_duration_ns: storage_pause_duration_ns,
293 buffer_unconsumed_item_count,
294 buffer_unconsumed_row_count,
295 buffer_unconsumed_epoch_count,
296 buffer_unconsumed_min_epoch,
297 buffer_memory_bytes,
298 buffer_read_count,
299 buffer_read_size,
300 total_read_count,
301 total_read_size,
302 persistent_log_read_metrics: KvLogStoreReadMetrics {
303 storage_read_size: persistent_log_read_size,
304 storage_read_count: persistent_log_read_count,
305 },
306 flushed_buffer_read_metrics: KvLogStoreReadMetrics {
307 storage_read_count: flushed_buffer_read_count,
308 storage_read_size: flushed_buffer_read_size,
309 },
310 }
311 }
312
313 #[cfg(test)]
314 pub(crate) fn for_test() -> Self {
315 SyncedKvLogStoreMetrics {
316 unclean_state: LabelGuardedIntCounter::test_int_counter::<5>(),
317 clean_state: LabelGuardedIntCounter::test_int_counter::<5>(),
318 wait_next_poll_ns: LabelGuardedIntCounter::test_int_counter::<4>(),
319 storage_write_count: LabelGuardedIntCounter::test_int_counter::<4>(),
320 storage_write_size: LabelGuardedIntCounter::test_int_counter::<4>(),
321 pause_duration_ns: LabelGuardedIntCounter::test_int_counter::<4>(),
322 buffer_unconsumed_item_count: LabelGuardedIntGauge::test_int_gauge::<4>(),
323 buffer_unconsumed_row_count: LabelGuardedIntGauge::test_int_gauge::<4>(),
324 buffer_unconsumed_epoch_count: LabelGuardedIntGauge::test_int_gauge::<4>(),
325 buffer_unconsumed_min_epoch: LabelGuardedIntGauge::test_int_gauge::<4>(),
326 buffer_memory_bytes: LabelGuardedIntGauge::test_int_gauge::<4>(),
327 buffer_read_count: LabelGuardedIntCounter::test_int_counter::<5>(),
328 buffer_read_size: LabelGuardedIntCounter::test_int_counter::<5>(),
329 total_read_count: LabelGuardedIntCounter::test_int_counter::<5>(),
330 total_read_size: LabelGuardedIntCounter::test_int_counter::<5>(),
331 persistent_log_read_metrics: KvLogStoreReadMetrics::for_test(),
332 flushed_buffer_read_metrics: KvLogStoreReadMetrics::for_test(),
333 }
334 }
335 }
336}
337
338pub(crate) type ReadFlushedChunkFuture =
339 BoxFuture<'static, LogStoreResult<(ChunkId, StreamChunk, Epoch)>>;
340
341pub struct SyncKvLogStoreContext<S: StateStore> {
342 pub table_id: TableId,
343 pub fragment_id: FragmentId,
344 pub metrics: SyncedKvLogStoreMetrics,
345 pub serde: LogStoreRowSerde,
346 pub state_store: S,
347 pub max_buffer_size: usize,
348 pub chunk_size: usize,
349 pub pause_duration_ms: Duration,
350 pub aligned: bool,
351}
352
353pub struct SyncedKvLogStoreExecutor<S: StateStore> {
354 actor_context: ActorContextRef,
355 upstream: Executor,
356 logstore_context: SyncKvLogStoreContext<S>,
357}
358
359impl<S: StateStore> SyncedKvLogStoreExecutor<S> {
361 #[expect(clippy::too_many_arguments)]
362 pub(crate) fn new(
363 actor_context: ActorContextRef,
364 table_id: TableId,
365 metrics: SyncedKvLogStoreMetrics,
366 serde: LogStoreRowSerde,
367 state_store: S,
368 buffer_size: usize,
369 chunk_size: usize,
370 upstream: Executor,
371 pause_duration_ms: Duration,
372 aligned: bool,
373 ) -> Self {
374 let logstore_context = SyncKvLogStoreContext {
375 table_id,
376 fragment_id: actor_context.fragment_id,
377 metrics,
378 serde,
379 state_store,
380 max_buffer_size: buffer_size,
381 chunk_size,
382 pause_duration_ms,
383 aligned,
384 };
385 Self {
386 actor_context,
387 upstream,
388 logstore_context,
389 }
390 }
391}
392
393pub(crate) struct FlushedChunkInfo {
394 epoch: u64,
395 start_seq_id: SeqId,
396 end_seq_id: SeqId,
397 flush_info: FlushInfo,
398 vnode_bitmap: Bitmap,
399}
400
401pub(crate) enum WriteFuture<S: LocalStateStore> {
402 Paused {
412 start_instant: Instant,
413 sleep_future: Option<Pin<Box<Sleep>>>,
414 barrier: Barrier,
415 stream: BoxedMessageStream,
416 write_state: LogStoreWriteState<S>, },
418 ReceiveFromUpstream {
419 future: StreamFuture<BoxedMessageStream>,
420 write_state: LogStoreWriteState<S>,
421 },
422 FlushingChunk {
423 epoch: u64,
424 start_seq_id: SeqId,
425 end_seq_id: SeqId,
426 future: Pin<Box<LogStoreStateWriteChunkFuture<S>>>,
427 stream: BoxedMessageStream,
428 },
429 Empty,
430}
431
432pub(crate) enum WriteFutureEvent {
433 UpstreamMessageReceived(Message),
434 ChunkFlushed(FlushedChunkInfo),
435}
436
437impl<S: LocalStateStore> WriteFuture<S> {
438 fn flush_chunk(
439 stream: BoxedMessageStream,
440 write_state: LogStoreWriteState<S>,
441 chunk: StreamChunk,
442 epoch: u64,
443 start_seq_id: SeqId,
444 end_seq_id: SeqId,
445 ) -> Self {
446 tracing::trace!(
447 start_seq_id,
448 end_seq_id,
449 epoch,
450 cardinality = chunk.cardinality(),
451 "write_future: flushing chunk"
452 );
453 Self::FlushingChunk {
454 epoch,
455 start_seq_id,
456 end_seq_id,
457 future: Box::pin(write_state.into_write_chunk_future(
458 chunk,
459 epoch,
460 start_seq_id,
461 end_seq_id,
462 )),
463 stream,
464 }
465 }
466
467 pub(crate) fn receive_from_upstream(
468 stream: BoxedMessageStream,
469 write_state: LogStoreWriteState<S>,
470 ) -> Self {
471 Self::ReceiveFromUpstream {
472 future: stream.into_future(),
473 write_state,
474 }
475 }
476
477 pub(crate) fn paused(
478 duration: Duration,
479 barrier: Barrier,
480 stream: BoxedMessageStream,
481 write_state: LogStoreWriteState<S>,
482 ) -> Self {
483 let now = Instant::now();
484 tracing::trace!(?now, ?duration, "write_future_pause");
485 Self::Paused {
486 start_instant: now,
487 sleep_future: Some(Box::pin(sleep_until(now + duration))),
488 barrier,
489 stream,
490 write_state,
491 }
492 }
493
494 pub(crate) async fn next_event(
495 &mut self,
496 metrics: &SyncedKvLogStoreMetrics,
497 ) -> StreamExecutorResult<(BoxedMessageStream, LogStoreWriteState<S>, WriteFutureEvent)> {
498 match self {
499 WriteFuture::Paused {
500 start_instant,
501 sleep_future,
502 ..
503 } => {
504 if let Some(sleep_future) = sleep_future {
505 sleep_future.await;
506 metrics
507 .pause_duration_ns
508 .inc_by(start_instant.elapsed().as_nanos() as _);
509 tracing::trace!("resuming write future");
510 }
511 must_match!(replace(self, WriteFuture::Empty), WriteFuture::Paused { stream, write_state, barrier, .. } => {
512 Ok((stream, write_state, WriteFutureEvent::UpstreamMessageReceived(Message::Barrier(barrier))))
513 })
514 }
515 WriteFuture::ReceiveFromUpstream { future, .. } => {
516 let (opt, stream) = future.await;
517 must_match!(replace(self, WriteFuture::Empty), WriteFuture::ReceiveFromUpstream { write_state, .. } => {
518 opt
519 .ok_or_else(|| anyhow!("end of upstream input").into())
520 .and_then(|result| result.map(|item| {
521 (stream, write_state, WriteFutureEvent::UpstreamMessageReceived(item))
522 }))
523 })
524 }
525 WriteFuture::FlushingChunk { future, .. } => {
526 let (write_state, result) = future.await;
527 let result = must_match!(replace(self, WriteFuture::Empty), WriteFuture::FlushingChunk { epoch, start_seq_id, end_seq_id, stream, .. } => {
528 result.map(|(flush_info, vnode_bitmap)| {
529 (stream, write_state, WriteFutureEvent::ChunkFlushed(FlushedChunkInfo {
530 epoch,
531 start_seq_id,
532 end_seq_id,
533 flush_info,
534 vnode_bitmap,
535 }))
536 })
537 });
538 result.map_err(Into::into)
539 }
540 WriteFuture::Empty => {
541 unreachable!("should not be polled after ready")
542 }
543 }
544 }
545}
546
547pub(crate) type LocalLogStoreReadState<S> =
548 LogStoreReadState<<<S as StateStore>::Local as LocalStateStore>::FlushedSnapshotReader>;
549pub(crate) type LocalLogStoreWriteState<S> = LogStoreWriteState<<S as StateStore>::Local>;
550
551impl<S: StateStore> SyncedKvLogStoreExecutor<S> {
553 pub(crate) async fn init_local_log_store_state(
554 context: &SyncKvLogStoreContext<S>,
555 first_write_epoch: EpochPair,
556 ) -> StreamExecutorResult<(LocalLogStoreReadState<S>, LocalLogStoreWriteState<S>)> {
557 let local_state_store = context
558 .state_store
559 .new_local(NewLocalOptions {
560 table_id: context.table_id,
561 fragment_id: context.fragment_id,
562 op_consistency_level: OpConsistencyLevel::Inconsistent,
563 table_option: TableOption {
564 retention_seconds: None,
565 },
566 is_replicated: false,
567 vnodes: context.serde.vnodes().clone(),
568 upload_on_flush: false,
569 })
570 .await;
571
572 let (read_state, mut initial_write_state) = new_log_store_state(
573 context.table_id,
574 local_state_store,
575 context.serde.clone(),
576 context.chunk_size,
577 );
578 initial_write_state.init(first_write_epoch).await?;
579 Ok((read_state, initial_write_state))
580 }
581
582 #[try_stream(ok = Message, error = StreamExecutorError)]
583 pub(crate) async fn aligned_message_stream(
584 actor_id: ActorId,
585 input: BoxedMessageStream,
586 read_state: LocalLogStoreReadState<S>,
587 mut initial_write_state: LocalLogStoreWriteState<S>,
588 metrics: SyncedKvLogStoreMetrics,
589 initial_write_epoch: EpochPair,
590 ) {
591 tracing::info!("aligned mode");
592 let log_store_stream = read_state
597 .read_persisted_log_store(
598 metrics.persistent_log_read_metrics.clone(),
599 initial_write_epoch.curr,
600 LogStoreReadStateStreamRangeStart::Unbounded,
601 )
602 .await?;
603
604 #[for_await]
605 for message in log_store_stream {
606 let (_epoch, message) = message?;
607 match message {
608 KvLogStoreItem::Barrier { .. } => {
609 continue;
610 }
611 KvLogStoreItem::StreamChunk { chunk, .. } => {
612 yield Message::Chunk(chunk);
613 }
614 }
615 }
616
617 let mut realigned_logstore = false;
618
619 #[for_await]
620 for message in input {
621 match message? {
622 Message::Barrier(barrier) => {
623 let is_checkpoint = barrier.is_checkpoint();
624 let mut progress = LogStoreVnodeProgress::None;
625 progress.apply_aligned(read_state.vnodes().clone(), barrier.epoch.prev, None);
626 let post_seal =
628 initial_write_state.seal_current_epoch(barrier.epoch.curr, progress.take());
629 barrier.assume_no_update_vnode_bitmap(actor_id)?;
630 yield Message::Barrier(barrier);
631 post_seal.post_yield_barrier(None).await?;
632 if !realigned_logstore && is_checkpoint {
633 realigned_logstore = true;
634 tracing::info!("realigned logstore");
635 }
636 }
637 Message::Chunk(chunk) => {
638 yield Message::Chunk(chunk);
639 }
640 Message::Watermark(watermark) => {
641 yield Message::Watermark(watermark);
642 }
643 }
644 }
645 }
646
647 pub(crate) fn apply_pause_resume_mutation(barrier: &Barrier, pause_stream: &mut bool) {
648 if let Some(mutation) = barrier.mutation.as_deref() {
649 match mutation {
650 Mutation::Pause => {
651 *pause_stream = true;
652 }
653 Mutation::Resume => {
654 *pause_stream = false;
655 }
656 _ => {}
657 }
658 }
659 }
660
661 pub(crate) fn process_upstream_chunk(
662 seq_id: SeqId,
663 stream: BoxedMessageStream,
664 write_state: LogStoreWriteState<S::Local>,
665 chunk: StreamChunk,
666 buffer: &mut SyncedLogStoreBuffer,
667 ) -> (SeqId, WriteFuture<S::Local>) {
668 let cardinality = chunk.cardinality();
669 if cardinality == 0 {
670 tracing::warn!(
671 epoch = write_state.epoch().curr,
672 "received empty chunk (cardinality=0), skipping"
673 );
674 return (
675 seq_id,
676 WriteFuture::receive_from_upstream(stream, write_state),
677 );
678 }
679
680 let start_seq_id = seq_id;
681 let new_seq_id = seq_id + cardinality as SeqId;
682 let end_seq_id = new_seq_id - 1;
683 let epoch = write_state.epoch().curr;
684 tracing::trace!(
685 start_seq_id,
686 end_seq_id,
687 new_seq_id,
688 epoch,
689 cardinality,
690 "received chunk"
691 );
692 let next_write_future = if let Some(chunk_to_flush) =
693 buffer.add_or_flush_chunk(start_seq_id, end_seq_id, chunk, epoch)
694 {
695 WriteFuture::flush_chunk(
696 stream,
697 write_state,
698 chunk_to_flush,
699 epoch,
700 start_seq_id,
701 end_seq_id,
702 )
703 } else {
704 WriteFuture::receive_from_upstream(stream, write_state)
705 };
706 (new_seq_id, next_write_future)
707 }
708
709 pub(crate) fn process_flushed_chunk(
710 stream: BoxedMessageStream,
711 write_state: LogStoreWriteState<S::Local>,
712 info: FlushedChunkInfo,
713 buffer: &mut SyncedLogStoreBuffer,
714 metrics: &SyncedKvLogStoreMetrics,
715 ) -> WriteFuture<S::Local> {
716 buffer.add_flushed_item_to_buffer(
717 info.start_seq_id,
718 info.end_seq_id,
719 info.vnode_bitmap,
720 info.epoch,
721 );
722 metrics
723 .storage_write_count
724 .inc_by(info.flush_info.flush_count as _);
725 metrics
726 .storage_write_size
727 .inc_by(info.flush_info.flush_size as _);
728 WriteFuture::receive_from_upstream(stream, write_state)
729 }
730
731 #[try_stream(ok= Message, error = StreamExecutorError)]
732 pub async fn execute_monitored(self) {
733 let wait_next_poll_ns = self.logstore_context.metrics.wait_next_poll_ns.clone();
734 #[for_await]
735 for message in self.execute_inner() {
736 let current_time = Instant::now();
737 yield message?;
738 wait_next_poll_ns.inc_by(current_time.elapsed().as_nanos() as _);
739 }
740 }
741
742 #[try_stream(ok = Message, error = StreamExecutorError)]
743 async fn execute_inner(self) {
744 let mut input = self.upstream.execute();
745
746 let first_barrier = expect_first_barrier(&mut input).await?;
748 let first_write_epoch = first_barrier.epoch;
749 yield Message::Barrier(first_barrier.clone());
750
751 let (read_state, initial_write_state) =
752 Self::init_local_log_store_state(&self.logstore_context, first_write_epoch).await?;
753
754 let mut pause_stream = first_barrier.is_pause_on_startup();
755 let initial_write_epoch = first_write_epoch;
756
757 if self.logstore_context.aligned {
758 let aligned_stream = Self::aligned_message_stream(
759 self.actor_context.id,
760 input,
761 read_state,
762 initial_write_state,
763 self.logstore_context.metrics.clone(),
764 initial_write_epoch,
765 );
766 #[for_await]
767 for message in aligned_stream {
768 yield message?;
769 }
770 return Ok(());
771 }
772
773 let mut seq_id = FIRST_SEQ_ID;
774 let mut buffer = SyncedLogStoreBuffer::new(
775 self.logstore_context.max_buffer_size,
776 self.logstore_context.chunk_size,
777 &self.logstore_context.metrics,
778 );
779
780 let log_store_stream = read_state
781 .read_persisted_log_store(
782 self.logstore_context
783 .metrics
784 .persistent_log_read_metrics
785 .clone(),
786 initial_write_epoch.curr,
787 LogStoreReadStateStreamRangeStart::Unbounded,
788 )
789 .await?;
790
791 let mut log_store_stream = tokio_stream::StreamExt::peekable(log_store_stream);
792 let mut clean_state = log_store_stream.peek().await.is_none();
793 tracing::trace!(?clean_state);
794
795 let mut read_future_state = ReadFuture::ReadingPersistedStream(log_store_stream);
796
797 let mut write_future_state = WriteFuture::receive_from_upstream(input, initial_write_state);
798
799 let mut progress = LogStoreVnodeProgress::None;
800
801 loop {
802 let select_result = {
803 let read_future = async {
804 if pause_stream {
805 pending().await
806 } else {
807 read_future_state
808 .next_message(
809 &mut progress,
810 &read_state,
811 &mut buffer,
812 &self.logstore_context.metrics,
813 )
814 .await
815 }
816 };
817 pin_mut!(read_future);
818 let write_future = write_future_state.next_event(&self.logstore_context.metrics);
819 pin_mut!(write_future);
820 let output = select(write_future, read_future).await;
821 drop_either_future(output)
822 };
823 match select_result {
824 Either::Left(result) => {
825 drop(write_future_state);
827 let (stream, mut write_state, either) = result?;
828 match either {
829 WriteFutureEvent::UpstreamMessageReceived(msg) => match msg {
830 Message::Barrier(barrier) => {
831 if clean_state && barrier.kind.is_checkpoint() && !buffer.is_empty()
832 {
833 write_future_state = WriteFuture::paused(
834 self.logstore_context.pause_duration_ms,
835 barrier,
836 stream,
837 write_state,
838 );
839 clean_state = false;
840 self.logstore_context.metrics.unclean_state.inc();
841 } else {
842 Self::apply_pause_resume_mutation(&barrier, &mut pause_stream);
843 let write_state_post_write_barrier = Self::write_barrier(
844 self.actor_context.id,
845 &mut write_state,
846 barrier.clone(),
847 &self.logstore_context.metrics,
848 progress.take(),
849 &mut buffer,
850 )
851 .await?;
852 seq_id = FIRST_SEQ_ID;
853 barrier.assume_no_update_vnode_bitmap(self.actor_context.id)?;
854
855 yield Message::Barrier(barrier);
856
857 write_state_post_write_barrier
858 .post_yield_barrier(None)
859 .await?;
860
861 write_future_state =
862 WriteFuture::receive_from_upstream(stream, write_state);
863 }
864 }
865 Message::Chunk(chunk) => {
866 let (new_seq_id, next_write_future) = Self::process_upstream_chunk(
867 seq_id,
868 stream,
869 write_state,
870 chunk,
871 &mut buffer,
872 );
873 seq_id = new_seq_id;
874 write_future_state = next_write_future;
875 }
876 Message::Watermark(watermark) => {
877 buffer.add_watermark(write_state.epoch().curr, watermark);
878 write_future_state =
879 WriteFuture::receive_from_upstream(stream, write_state);
880 }
881 },
882 WriteFutureEvent::ChunkFlushed(info) => {
883 write_future_state = Self::process_flushed_chunk(
884 stream,
885 write_state,
886 info,
887 &mut buffer,
888 &self.logstore_context.metrics,
889 );
890 }
891 }
892 }
893 Either::Right(result) => {
894 let clean_state_reached = read_future_state.mark_clean_state(
895 &mut clean_state,
896 &buffer,
897 &self.logstore_context.metrics,
898 );
899
900 if clean_state_reached {
901 if let WriteFuture::Paused { sleep_future, .. } = &mut write_future_state {
903 tracing::trace!("resuming paused future");
904 assert!(buffer.has_available_capacity());
905 *sleep_future = None;
906 }
907 }
908 let message = result?;
909 if let Message::Chunk(chunk) = &message {
910 self.logstore_context
911 .metrics
912 .total_read_count
913 .inc_by(chunk.cardinality() as _);
914 }
915
916 yield message;
917 }
918 }
919 }
920 }
921}
922
923pub(crate) type PersistedStream<S> =
924 Peekable<Pin<Box<LogStoreItemMergeStream<TimeoutAutoRebuildIter<S>>>>>;
925
926pub(crate) enum ReadFuture<S: StateStoreRead> {
927 ReadingPersistedStream(PersistedStream<S>),
928 ReadingFlushedChunk {
929 future: ReadFlushedChunkFuture,
930 end_seq_id: SeqId,
931 },
932 Idle,
933}
934
935impl<S: StateStoreRead> ReadFuture<S> {
937 pub(crate) fn mark_clean_state(
938 &self,
939 clean_state: &mut bool,
940 buffer: &SyncedLogStoreBuffer,
941 metrics: &SyncedKvLogStoreMetrics,
942 ) -> bool {
943 if !*clean_state && matches!(self, ReadFuture::Idle) && buffer.is_empty() {
944 *clean_state = true;
945 metrics.clean_state.inc();
946 true
947 } else {
948 false
949 }
950 }
951
952 pub(crate) async fn next_message(
953 &mut self,
954 progress: &mut LogStoreVnodeProgress,
955 read_state: &LogStoreReadState<S>,
956 buffer: &mut SyncedLogStoreBuffer,
957 metrics: &SyncedKvLogStoreMetrics,
958 ) -> StreamExecutorResult<Message> {
959 match self {
960 ReadFuture::ReadingPersistedStream(stream) => {
961 while let Some((epoch, item)) = stream.try_next().await? {
962 match item {
963 KvLogStoreItem::Barrier { vnodes, .. } => {
964 tracing::trace!(epoch, "read logstore barrier");
965 progress.apply_aligned(vnodes, epoch, None);
967 continue;
968 }
969 KvLogStoreItem::StreamChunk {
970 chunk,
971 progress: chunk_progress,
972 } => {
973 tracing::trace!("read logstore chunk of size: {}", chunk.cardinality());
974 progress.apply_per_vnode(epoch, chunk_progress);
975 return Ok(Message::Chunk(chunk));
976 }
977 }
978 }
979 *self = ReadFuture::Idle;
980 }
981 ReadFuture::ReadingFlushedChunk { .. } | ReadFuture::Idle => {}
982 }
983 match self {
984 ReadFuture::ReadingPersistedStream(_) => {
985 unreachable!("must have finished read persisted stream when reaching here")
986 }
987 ReadFuture::ReadingFlushedChunk { .. } => {}
988 ReadFuture::Idle => loop {
989 let Some((item_epoch, item)) = buffer.pop_front() else {
990 return pending().await;
991 };
992 match item {
993 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
994 chunk,
995 start_seq_id,
996 end_seq_id,
997 flushed,
998 ..
999 }) => {
1000 metrics.buffer_read_count.inc_by(chunk.cardinality() as _);
1001 tracing::trace!(
1002 start_seq_id,
1003 end_seq_id,
1004 flushed,
1005 cardinality = chunk.cardinality(),
1006 "read buffered chunk of size"
1007 );
1008 progress.apply_aligned(
1009 read_state.vnodes().clone(),
1010 item_epoch,
1011 Some(end_seq_id),
1012 );
1013 return Ok(Message::Chunk(chunk));
1014 }
1015 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1016 vnode_bitmap,
1017 start_seq_id,
1018 end_seq_id,
1019 chunk_id,
1020 }) => {
1021 tracing::trace!(start_seq_id, end_seq_id, chunk_id, "read flushed chunk");
1022 let read_metrics = metrics.flushed_buffer_read_metrics.clone();
1023 let future = read_state
1024 .read_flushed_chunk(
1025 vnode_bitmap,
1026 chunk_id,
1027 start_seq_id,
1028 end_seq_id,
1029 item_epoch,
1030 read_metrics,
1031 )
1032 .boxed();
1033 *self = ReadFuture::ReadingFlushedChunk { future, end_seq_id };
1034 break;
1035 }
1036 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Barrier { .. }) => {
1037 tracing::trace!(item_epoch, "read buffer barrier");
1038 progress.apply_aligned(read_state.vnodes().clone(), item_epoch, None);
1039 continue;
1040 }
1041 SyncedLogStoreBufferItem::Watermark(watermark) => {
1042 return Ok(Message::Watermark(watermark));
1043 }
1044 }
1045 },
1046 }
1047
1048 let (future, end_seq_id) = match self {
1049 ReadFuture::ReadingPersistedStream(_) | ReadFuture::Idle => {
1050 unreachable!("should be at ReadingFlushedChunk")
1051 }
1052 ReadFuture::ReadingFlushedChunk { future, end_seq_id } => (future, *end_seq_id),
1053 };
1054
1055 let (_, chunk, epoch) = future.await?;
1056 progress.apply_aligned(read_state.vnodes().clone(), epoch, Some(end_seq_id));
1057 tracing::trace!(
1058 end_seq_id,
1059 "read flushed chunk of size: {}",
1060 chunk.cardinality()
1061 );
1062 *self = ReadFuture::Idle;
1063 Ok(Message::Chunk(chunk))
1064 }
1065}
1066
1067impl<S: StateStore> SyncedKvLogStoreExecutor<S> {
1069 pub(crate) async fn write_barrier<'a>(
1070 actor_id: ActorId,
1071 write_state: &'a mut LogStoreWriteState<S::Local>,
1072 barrier: Barrier,
1073 metrics: &SyncedKvLogStoreMetrics,
1074 progress: LogStoreVnodeProgress,
1075 buffer: &mut SyncedLogStoreBuffer,
1076 ) -> StreamExecutorResult<LogStorePostSealCurrentEpoch<'a, S::Local>> {
1077 tracing::trace!(%actor_id, ?progress, "applying truncation");
1078 let epoch = barrier.epoch.prev;
1082 let mut writer = write_state.start_writer(false);
1083 writer.write_barrier(epoch, barrier.is_checkpoint())?;
1084
1085 if barrier.is_checkpoint() {
1086 for (epoch, item) in buffer.buffer.iter_mut().rev() {
1087 match item {
1088 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1089 chunk,
1090 start_seq_id,
1091 end_seq_id,
1092 flushed,
1093 ..
1094 }) => {
1095 if !*flushed {
1096 writer.write_chunk(chunk, *epoch, *start_seq_id, *end_seq_id)?;
1097 *flushed = true;
1098 } else {
1099 break;
1100 }
1101 }
1102 SyncedLogStoreBufferItem::LogStore(
1103 LogStoreBufferItem::Flushed { .. } | LogStoreBufferItem::Barrier { .. },
1104 )
1105 | SyncedLogStoreBufferItem::Watermark(_) => {}
1106 }
1107 }
1108 }
1109
1110 let (flush_info, _) = writer.finish().await?;
1112 metrics
1113 .storage_write_count
1114 .inc_by(flush_info.flush_count as _);
1115 metrics
1116 .storage_write_size
1117 .inc_by(flush_info.flush_size as _);
1118 let post_seal = write_state.seal_current_epoch(barrier.epoch.curr, progress);
1119
1120 buffer.buffer.push_back((
1122 epoch,
1123 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Barrier {
1124 is_checkpoint: barrier.is_checkpoint(),
1125 next_epoch: barrier.epoch.curr,
1126 schema_change: None,
1127 is_stop: false,
1128 }),
1129 ));
1130 buffer.next_chunk_id = 0;
1131 buffer.update_buffer_metrics();
1132
1133 Ok(post_seal)
1134 }
1135}
1136
1137#[derive(EstimateSize)]
1138enum SyncedLogStoreBufferItem {
1139 LogStore(LogStoreBufferItem),
1140 Watermark(Watermark),
1141}
1142
1143pub(crate) struct SyncedLogStoreBuffer {
1144 buffer: VecDeque<(u64, SyncedLogStoreBufferItem)>,
1145 current_size: usize,
1146 max_size: usize,
1147 max_chunk_size: usize,
1148 next_chunk_id: ChunkId,
1149 metrics: SyncedKvLogStoreMetrics,
1150 flushed_count: usize,
1151}
1152
1153impl SyncedLogStoreBuffer {
1154 pub fn new(
1155 max_buffer_size: usize,
1156 chunk_size: usize,
1157 metrics: &SyncedKvLogStoreMetrics,
1158 ) -> Self {
1159 SyncedLogStoreBuffer {
1160 buffer: VecDeque::new(),
1161 current_size: 0,
1162 max_size: max_buffer_size,
1163 max_chunk_size: chunk_size,
1164 next_chunk_id: 0,
1165 metrics: metrics.clone(),
1166 flushed_count: 0,
1167 }
1168 }
1169
1170 pub fn is_empty(&self) -> bool {
1171 self.current_size == 0
1172 }
1173
1174 pub fn has_available_capacity(&self) -> bool {
1175 self.current_size < self.max_size
1176 }
1177
1178 fn add_or_flush_chunk(
1179 &mut self,
1180 start_seq_id: SeqId,
1181 end_seq_id: SeqId,
1182 chunk: StreamChunk,
1183 epoch: u64,
1184 ) -> Option<StreamChunk> {
1185 let current_size = self.current_size;
1186 let chunk_size = chunk.cardinality();
1187
1188 tracing::trace!(
1189 current_size,
1190 chunk_size,
1191 max_size = self.max_size,
1192 "checking chunk size"
1193 );
1194 let should_flush_chunk = current_size + chunk_size > self.max_size;
1195 if should_flush_chunk {
1196 tracing::trace!(start_seq_id, end_seq_id, epoch, "flushing chunk",);
1197 Some(chunk)
1198 } else {
1199 tracing::trace!(start_seq_id, end_seq_id, epoch, "buffering chunk",);
1200 self.add_chunk_to_buffer(chunk, start_seq_id, end_seq_id, epoch);
1201 None
1202 }
1203 }
1204
1205 fn add_flushed_item_to_buffer(
1208 &mut self,
1209 start_seq_id: SeqId,
1210 end_seq_id: SeqId,
1211 new_vnode_bitmap: Bitmap,
1212 epoch: u64,
1213 ) {
1214 let new_chunk_size = (end_seq_id - start_seq_id + 1) as usize;
1215
1216 if let Some((
1217 item_epoch,
1218 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1219 start_seq_id: prev_start_seq_id,
1220 end_seq_id: prev_end_seq_id,
1221 vnode_bitmap,
1222 ..
1223 }),
1224 )) = self.buffer.back_mut()
1225 && let flushed_chunk_size = (*prev_end_seq_id - *prev_start_seq_id + 1) as usize
1226 && let projected_flushed_chunk_size = flushed_chunk_size + new_chunk_size
1227 && projected_flushed_chunk_size <= self.max_chunk_size
1228 {
1229 assert!(
1230 *prev_end_seq_id < start_seq_id,
1231 "prev end_seq_id {} should be smaller than current start_seq_id {}",
1232 end_seq_id,
1233 start_seq_id
1234 );
1235 assert_eq!(
1236 epoch, *item_epoch,
1237 "epoch of newly added flushed item must be the same as the last flushed item"
1238 );
1239 *prev_end_seq_id = end_seq_id;
1240 *vnode_bitmap |= new_vnode_bitmap;
1241 } else {
1242 let chunk_id = self.next_chunk_id;
1243 self.next_chunk_id += 1;
1244 self.buffer.push_back((
1245 epoch,
1246 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1247 start_seq_id,
1248 end_seq_id,
1249 vnode_bitmap: new_vnode_bitmap,
1250 chunk_id,
1251 }),
1252 ));
1253 self.flushed_count += 1;
1254 tracing::trace!(
1255 "adding flushed item to buffer: start_seq_id: {start_seq_id}, end_seq_id: {end_seq_id}, chunk_id: {chunk_id}"
1256 );
1257 }
1258 self.update_buffer_metrics();
1260 }
1261
1262 fn add_chunk_to_buffer(
1263 &mut self,
1264 chunk: StreamChunk,
1265 start_seq_id: SeqId,
1266 end_seq_id: SeqId,
1267 epoch: u64,
1268 ) {
1269 let chunk_id = self.next_chunk_id;
1270 self.next_chunk_id += 1;
1271 self.current_size += chunk.cardinality();
1272 self.buffer.push_back((
1273 epoch,
1274 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1275 chunk,
1276 start_seq_id,
1277 end_seq_id,
1278 flushed: false,
1279 chunk_id,
1280 }),
1281 ));
1282 self.update_buffer_metrics();
1283 }
1284
1285 pub(crate) fn add_watermark(&mut self, epoch: u64, watermark: Watermark) {
1286 self.buffer
1287 .push_back((epoch, SyncedLogStoreBufferItem::Watermark(watermark)));
1288 self.update_buffer_metrics();
1289 }
1290
1291 fn pop_front(&mut self) -> Option<(u64, SyncedLogStoreBufferItem)> {
1292 let item = self.buffer.pop_front();
1293 match &item {
1294 Some((_, SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed { .. }))) => {
1295 self.flushed_count -= 1;
1296 }
1297 Some((
1298 _,
1299 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1300 chunk, ..
1301 }),
1302 )) => {
1303 self.current_size -= chunk.cardinality();
1304 }
1305 _ => {}
1306 }
1307 self.update_buffer_metrics();
1308 item
1309 }
1310
1311 fn update_buffer_metrics(&self) {
1312 let mut epoch_count = 0;
1313 let mut row_count = 0;
1314 let mut memory_bytes = 0;
1315 for (_, item) in &self.buffer {
1316 match item {
1317 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::StreamChunk {
1318 chunk,
1319 ..
1320 }) => {
1321 row_count += chunk.cardinality();
1322 }
1323 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Flushed {
1324 start_seq_id,
1325 end_seq_id,
1326 ..
1327 }) => {
1328 row_count += (end_seq_id - start_seq_id) as usize;
1329 }
1330 SyncedLogStoreBufferItem::LogStore(LogStoreBufferItem::Barrier { .. }) => {
1331 epoch_count += 1;
1332 }
1333 SyncedLogStoreBufferItem::Watermark(_) => {}
1334 }
1335 memory_bytes += item.estimated_size();
1336 }
1337 self.metrics.buffer_unconsumed_epoch_count.set(epoch_count);
1338 self.metrics.buffer_unconsumed_row_count.set(row_count as _);
1339 self.metrics
1340 .buffer_unconsumed_item_count
1341 .set(self.buffer.len() as _);
1342 self.metrics.buffer_unconsumed_min_epoch.set(
1343 self.buffer
1344 .front()
1345 .map(|(epoch, _)| *epoch)
1346 .unwrap_or_default() as _,
1347 );
1348 self.metrics.buffer_memory_bytes.set(memory_bytes as _);
1349 }
1350}
1351
1352impl<S> Execute for SyncedKvLogStoreExecutor<S>
1353where
1354 S: StateStore,
1355{
1356 fn execute(self: Box<Self>) -> BoxedMessageStream {
1357 self.execute_monitored().boxed()
1358 }
1359}
1360
1361#[cfg(test)]
1362mod tests {
1363 use itertools::Itertools;
1364 use pretty_assertions::assert_eq;
1365 use risingwave_common::catalog::Field;
1366 use risingwave_common::hash::VirtualNode;
1367 use risingwave_common::test_prelude::*;
1368 use risingwave_common::util::epoch::test_epoch;
1369 use risingwave_storage::memory::MemoryStateStore;
1370
1371 use super::*;
1372 use crate::assert_stream_chunk_eq;
1373 use crate::common::log_store_impl::kv_log_store::KV_LOG_STORE_V2_INFO;
1374 use crate::common::log_store_impl::kv_log_store::test_utils::{
1375 check_stream_chunk_eq, gen_test_log_store_table, test_payload_schema,
1376 };
1377 use crate::executor::sync_kv_log_store::metrics::SyncedKvLogStoreMetrics;
1378 use crate::executor::test_utils::MockSource;
1379
1380 fn init_logger() {
1381 let _ = tracing_subscriber::fmt()
1382 .with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
1383 .with_ansi(false)
1384 .try_init();
1385 }
1386
1387 #[tokio::test]
1389 async fn test_read_write_buffer() {
1390 init_logger();
1391
1392 let pk_info = &KV_LOG_STORE_V2_INFO;
1393 let column_descs = test_payload_schema(pk_info);
1394 let fields = column_descs
1395 .into_iter()
1396 .map(|desc| Field::new(desc.name.clone(), desc.data_type))
1397 .collect_vec();
1398 let schema = Schema { fields };
1399 let stream_key = vec![0];
1400 let (mut tx, source) = MockSource::channel();
1401 let source = source.into_executor(schema.clone(), stream_key.clone());
1402
1403 let vnodes = Some(Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)));
1404
1405 let table = gen_test_log_store_table(pk_info);
1406
1407 let log_store_executor = SyncedKvLogStoreExecutor::new(
1408 ActorContext::for_test(123),
1409 table.id,
1410 SyncedKvLogStoreMetrics::for_test(),
1411 LogStoreRowSerde::new(&table, vnodes, pk_info),
1412 MemoryStateStore::new(),
1413 10,
1414 256,
1415 source,
1416 Duration::from_millis(256),
1417 false,
1418 )
1419 .boxed();
1420
1421 tx.push_barrier(test_epoch(1), false);
1423
1424 let chunk_1 = StreamChunk::from_pretty(
1425 " I T
1426 + 5 10
1427 + 6 10
1428 + 8 10
1429 + 9 10
1430 + 10 11",
1431 );
1432
1433 let chunk_2 = StreamChunk::from_pretty(
1434 " I T
1435 - 5 10
1436 - 6 10
1437 - 8 10
1438 U- 9 10
1439 U+ 10 11",
1440 );
1441
1442 tx.push_chunk(chunk_1.clone());
1443 tx.push_int64_watermark(0, 7);
1444 tx.push_chunk(chunk_2.clone());
1445
1446 let mut stream = log_store_executor.execute();
1447
1448 match stream.next().await {
1449 Some(Ok(Message::Barrier(barrier))) => {
1450 assert_eq!(barrier.epoch.curr, test_epoch(1));
1451 }
1452 other => panic!("Expected a barrier message, got {:?}", other),
1453 }
1454
1455 match stream.next().await {
1456 Some(Ok(Message::Chunk(chunk))) => {
1457 assert_stream_chunk_eq!(chunk, chunk_1);
1458 }
1459 other => panic!("Expected a chunk message, got {:?}", other),
1460 }
1461
1462 match stream.next().await {
1463 Some(Ok(Message::Watermark(watermark))) => {
1464 assert_eq!(watermark, Watermark::new(0, DataType::Int64, 7_i64.into()));
1465 }
1466 other => panic!("Expected a watermark message, got {:?}", other),
1467 }
1468
1469 match stream.next().await {
1470 Some(Ok(Message::Chunk(chunk))) => {
1471 assert_stream_chunk_eq!(chunk, chunk_2);
1472 }
1473 other => panic!("Expected a chunk message, got {:?}", other),
1474 }
1475
1476 tx.push_barrier(test_epoch(2), false);
1477
1478 match stream.next().await {
1479 Some(Ok(Message::Barrier(barrier))) => {
1480 assert_eq!(barrier.epoch.curr, test_epoch(2));
1481 }
1482 other => panic!("Expected a barrier message, got {:?}", other),
1483 }
1484 }
1485
1486 #[tokio::test]
1492 async fn test_barrier_persisted_read() {
1493 init_logger();
1494
1495 let pk_info = &KV_LOG_STORE_V2_INFO;
1496 let column_descs = test_payload_schema(pk_info);
1497 let fields = column_descs
1498 .into_iter()
1499 .map(|desc| Field::new(desc.name.clone(), desc.data_type))
1500 .collect_vec();
1501 let schema = Schema { fields };
1502 let stream_key = vec![0];
1503 let (mut tx, source) = MockSource::channel();
1504 let source = source.into_executor(schema.clone(), stream_key.clone());
1505
1506 let vnodes = Some(Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)));
1507
1508 let table = gen_test_log_store_table(pk_info);
1509
1510 let log_store_executor = SyncedKvLogStoreExecutor::new(
1511 ActorContext::for_test(123),
1512 table.id,
1513 SyncedKvLogStoreMetrics::for_test(),
1514 LogStoreRowSerde::new(&table, vnodes, pk_info),
1515 MemoryStateStore::new(),
1516 10,
1517 256,
1518 source,
1519 Duration::from_millis(256),
1520 false,
1521 )
1522 .boxed();
1523
1524 tx.push_barrier(test_epoch(1), false);
1526
1527 let chunk_1 = StreamChunk::from_pretty(
1528 " I T
1529 + 5 10
1530 + 6 10
1531 + 8 10
1532 + 9 10
1533 + 10 11",
1534 );
1535
1536 let chunk_2 = StreamChunk::from_pretty(
1537 " I T
1538 - 5 10
1539 - 6 10
1540 - 8 10
1541 U- 10 11
1542 U+ 10 10",
1543 );
1544
1545 tx.push_chunk(chunk_1.clone());
1546 tx.push_int64_watermark(0, 7);
1547 tx.push_chunk(chunk_2.clone());
1548
1549 tx.push_barrier(test_epoch(2), false);
1550
1551 let mut stream = log_store_executor.execute();
1552
1553 match stream.next().await {
1554 Some(Ok(Message::Barrier(barrier))) => {
1555 assert_eq!(barrier.epoch.curr, test_epoch(1));
1556 }
1557 other => panic!("Expected a barrier message, got {:?}", other),
1558 }
1559
1560 match stream.next().await {
1561 Some(Ok(Message::Chunk(chunk))) => {
1562 assert_stream_chunk_eq!(chunk, chunk_1);
1563 }
1564 other => panic!("Expected a chunk message, got {:?}", other),
1565 }
1566
1567 match stream.next().await {
1568 Some(Ok(Message::Watermark(watermark))) => {
1569 assert_eq!(watermark, Watermark::new(0, DataType::Int64, 7_i64.into()));
1570 }
1571 other => panic!("Expected a watermark message, got {:?}", other),
1572 }
1573
1574 match stream.next().await {
1575 Some(Ok(Message::Chunk(chunk))) => {
1576 assert_stream_chunk_eq!(chunk, chunk_2);
1577 }
1578 other => panic!("Expected a chunk message, got {:?}", other),
1579 }
1580
1581 match stream.next().await {
1582 Some(Ok(Message::Barrier(barrier))) => {
1583 assert_eq!(barrier.epoch.curr, test_epoch(2));
1584 }
1585 other => panic!("Expected a barrier message, got {:?}", other),
1586 }
1587 }
1588
1589 #[tokio::test]
1592 async fn test_max_chunk_persisted_read() {
1593 init_logger();
1594
1595 let pk_info = &KV_LOG_STORE_V2_INFO;
1596 let column_descs = test_payload_schema(pk_info);
1597 let fields = column_descs
1598 .into_iter()
1599 .map(|desc| Field::new(desc.name.clone(), desc.data_type))
1600 .collect_vec();
1601 let schema = Schema { fields };
1602 let stream_key = vec![0];
1603 let (mut tx, source) = MockSource::channel();
1604 let source = source.into_executor(schema.clone(), stream_key.clone());
1605
1606 let vnodes = Some(Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)));
1607
1608 let table = gen_test_log_store_table(pk_info);
1609
1610 let log_store_executor = SyncedKvLogStoreExecutor::new(
1611 ActorContext::for_test(123),
1612 table.id,
1613 SyncedKvLogStoreMetrics::for_test(),
1614 LogStoreRowSerde::new(&table, vnodes, pk_info),
1615 MemoryStateStore::new(),
1616 0,
1617 256,
1618 source,
1619 Duration::from_millis(256),
1620 false,
1621 )
1622 .boxed();
1623
1624 tx.push_barrier(test_epoch(1), false);
1626
1627 let chunk_1 = StreamChunk::from_pretty(
1628 " I T
1629 + 5 10
1630 + 6 10
1631 + 8 10
1632 + 9 10
1633 + 10 11",
1634 );
1635
1636 let chunk_2 = StreamChunk::from_pretty(
1637 " I T
1638 - 5 10
1639 - 6 10
1640 - 8 10
1641 U- 10 11
1642 U+ 10 10",
1643 );
1644
1645 let chunk_3 = StreamChunk::from_pretty(
1646 " I T
1647 + 11 12",
1648 );
1649
1650 tx.push_chunk(chunk_1.clone());
1651 tx.push_chunk(chunk_2.clone());
1652 tx.push_int64_watermark(0, 7);
1653 tx.push_chunk(chunk_3.clone());
1654
1655 tx.push_barrier(test_epoch(2), false);
1656
1657 let mut stream = log_store_executor.execute();
1658
1659 for i in 1..=2 {
1660 match stream.next().await {
1661 Some(Ok(Message::Barrier(barrier))) => {
1662 assert_eq!(barrier.epoch.curr, test_epoch(i));
1663 }
1664 other => panic!("Expected a barrier message, got {:?}", other),
1665 }
1666 }
1667
1668 match stream.next().await {
1669 Some(Ok(Message::Chunk(actual))) => {
1670 let expected = StreamChunk::from_pretty(
1671 " I T
1672 + 5 10
1673 + 6 10
1674 + 8 10
1675 + 9 10
1676 + 10 11
1677 - 5 10
1678 - 6 10
1679 - 8 10
1680 U- 10 11
1681 U+ 10 10",
1682 );
1683 assert_stream_chunk_eq!(actual, expected);
1684 }
1685 other => panic!("Expected a chunk message, got {:?}", other),
1686 }
1687
1688 match stream.next().await {
1689 Some(Ok(Message::Watermark(watermark))) => {
1690 assert_eq!(watermark, Watermark::new(0, DataType::Int64, 7_i64.into()));
1691 }
1692 other => panic!("Expected a watermark message, got {:?}", other),
1693 }
1694
1695 match stream.next().await {
1696 Some(Ok(Message::Chunk(actual))) => {
1697 assert_stream_chunk_eq!(actual, chunk_3);
1698 }
1699 other => panic!("Expected a chunk message, got {:?}", other),
1700 }
1701 }
1702}