1use std::collections::{BTreeMap, HashSet, btree_map};
16use std::marker::PhantomData;
17use std::ops::RangeInclusive;
18
19use delta_btree_map::Change;
20use itertools::Itertools;
21use risingwave_common::array::stream_record::Record;
22use risingwave_common::config::streaming::OverWindowCachePolicy as CachePolicy;
23use risingwave_common::row::RowExt;
24use risingwave_common::types::{DefaultOrd, DefaultOrdered, ScalarImpl};
25use risingwave_common::util::memcmp_encoding::{self, MemcmpEncoded};
26use risingwave_common::util::sort_util::OrderType;
27use risingwave_expr::window_function::{
28 RangeFrameBounds, RowsFrameBounds, StateKey, WindowFuncCall, can_forward_watermark_on_order_key,
29};
30
31use super::frame_finder::merge_rows_frames;
32use super::over_partition::{OverPartition, PartitionDelta};
33use super::range_cache::{CacheKey, PartitionCache};
34use crate::cache::ManagedLruCache;
35use crate::common::change_buffer::ChangeBuffer;
36use crate::common::metrics::MetricsInfo;
37use crate::consistency::consistency_panic;
38use crate::executor::monitor::OverWindowMetrics;
39use crate::executor::prelude::*;
40
41pub struct OverWindowExecutor<S: StateStore> {
53 input: Executor,
54 inner: ExecutorInner<S>,
55}
56
57struct ExecutorInner<S: StateStore> {
58 actor_ctx: ActorContextRef,
59
60 schema: Schema,
61 calls: Calls,
62 deduped_part_key_indices: Vec<usize>,
63 order_key_indices: Vec<usize>,
64 order_key_data_types: Vec<DataType>,
65 order_key_order_types: Vec<OrderType>,
66 input_stream_key: Vec<usize>,
67 state_key_to_table_sub_pk_proj: Vec<usize>,
68
69 state_table: StateTable<S>,
70 watermark_sequence: AtomicU64Ref,
71
72 chunk_size: usize,
74 cache_policy: CachePolicy,
75 state_cleaning: Option<StateCleaning>,
77 watermark_forward_cols: HashSet<usize>,
80}
81
82struct ExecutionVars<S: StateStore> {
83 cached_partitions: ManagedLruCache<OwnedRow, PartitionCache>,
85 recently_accessed_ranges: BTreeMap<DefaultOrdered<OwnedRow>, RangeInclusive<StateKey>>,
87 cleaning_watermark: Option<ScalarImpl>,
89 touched_partitions: HashSet<OwnedRow>,
91 stats: ExecutionStats,
92 _phantom: PhantomData<S>,
93}
94
95#[derive(Debug)]
114pub(super) struct StateCleaning {
115 pub watermark_col_idx: usize,
117 pub stale_rows_at_front: bool,
120 pub n_retain: usize,
123}
124
125#[derive(Default)]
126struct ExecutionStats {
127 cache_miss: u64,
128 cache_lookup: u64,
129}
130
131impl<S: StateStore> Execute for OverWindowExecutor<S> {
132 fn execute(self: Box<Self>) -> crate::executor::BoxedMessageStream {
133 self.executor_inner().boxed()
134 }
135}
136
137impl<S: StateStore> ExecutorInner<S> {
138 fn get_partition_key(&self, full_row: impl Row) -> OwnedRow {
140 full_row
141 .project(&self.deduped_part_key_indices)
142 .into_owned_row()
143 }
144
145 fn get_input_pk(&self, full_row: impl Row) -> OwnedRow {
146 full_row.project(&self.input_stream_key).into_owned_row()
147 }
148
149 fn encode_order_key(&self, full_row: impl Row) -> StreamExecutorResult<MemcmpEncoded> {
151 Ok(memcmp_encoding::encode_row(
152 full_row.project(&self.order_key_indices),
153 &self.order_key_order_types,
154 )?)
155 }
156
157 fn row_to_cache_key(&self, full_row: impl Row + Copy) -> StreamExecutorResult<CacheKey> {
158 Ok(CacheKey::Normal(StateKey {
159 order_key: self.encode_order_key(full_row)?,
160 pk: self.get_input_pk(full_row).into(),
161 }))
162 }
163}
164
165pub struct OverWindowExecutorArgs<S: StateStore> {
166 pub actor_ctx: ActorContextRef,
167
168 pub input: Executor,
169
170 pub schema: Schema,
171 pub calls: Vec<WindowFuncCall>,
172 pub partition_key_indices: Vec<usize>,
173 pub order_key_indices: Vec<usize>,
174 pub order_key_order_types: Vec<OrderType>,
175
176 pub state_table: StateTable<S>,
177 pub watermark_epoch: AtomicU64Ref,
178 pub metrics: Arc<StreamingMetrics>,
179
180 pub chunk_size: usize,
181 pub cache_policy: CachePolicy,
182 pub enable_state_cleaning: bool,
184}
185
186pub(super) struct Calls {
190 calls: Vec<WindowFuncCall>,
191
192 pub(super) super_rows_frame_bounds: RowsFrameBounds,
194 pub(super) range_frames: Vec<RangeFrameBounds>,
196 pub(super) start_is_unbounded: bool,
197 pub(super) end_is_unbounded: bool,
198 pub(super) all_arg_indices: Vec<usize>,
200
201 pub(super) numbering_only: bool,
204 pub(super) has_rank: bool,
205}
206
207impl Calls {
208 fn new(calls: Vec<WindowFuncCall>) -> Self {
209 let rows_frames = calls
210 .iter()
211 .filter_map(|call| call.frame.bounds.as_rows())
212 .collect::<Vec<_>>();
213 let super_rows_frame_bounds = merge_rows_frames(&rows_frames);
214 let range_frames = calls
215 .iter()
216 .filter_map(|call| call.frame.bounds.as_range())
217 .cloned()
218 .collect::<Vec<_>>();
219
220 let start_is_unbounded = calls
221 .iter()
222 .any(|call| call.frame.bounds.start_is_unbounded());
223 let end_is_unbounded = calls
224 .iter()
225 .any(|call| call.frame.bounds.end_is_unbounded());
226
227 let all_arg_indices = calls
228 .iter()
229 .flat_map(|call| call.args.val_indices().iter().copied())
230 .dedup()
231 .collect();
232
233 let numbering_only = calls.iter().all(|call| call.kind.is_numbering());
234 let has_rank = calls.iter().any(|call| call.kind.is_rank());
235
236 Self {
237 calls,
238 super_rows_frame_bounds,
239 range_frames,
240 start_is_unbounded,
241 end_is_unbounded,
242 all_arg_indices,
243 numbering_only,
244 has_rank,
245 }
246 }
247
248 pub(super) fn iter(&self) -> impl ExactSizeIterator<Item = &WindowFuncCall> {
249 self.calls.iter()
250 }
251
252 pub(super) fn len(&self) -> usize {
253 self.calls.len()
254 }
255}
256
257impl<S: StateStore> OverWindowExecutor<S> {
258 pub fn new(args: OverWindowExecutorArgs<S>) -> Self {
259 let calls = Calls::new(args.calls);
260
261 let input_info = args.input.info().clone();
262 let input_schema = &input_info.schema;
263
264 let has_unbounded_frame = calls.start_is_unbounded || calls.end_is_unbounded;
265 let cache_policy = if has_unbounded_frame {
266 CachePolicy::Full
269 } else {
270 args.cache_policy
271 };
272
273 let order_key_data_types = args
274 .order_key_indices
275 .iter()
276 .map(|i| input_schema[*i].data_type())
277 .collect();
278
279 let state_key_to_table_sub_pk_proj = RowConverter::calc_state_key_to_table_sub_pk_proj(
280 &args.partition_key_indices,
281 &args.order_key_indices,
282 &input_info.stream_key,
283 );
284
285 let deduped_part_key_indices: Vec<usize> = {
286 let mut dedup = HashSet::new();
287 args.partition_key_indices
288 .iter()
289 .filter(|i| dedup.insert(**i))
290 .copied()
291 .collect()
292 };
293
294 let state_cleaning = if args.enable_state_cleaning {
295 let all_frames_bounded_rows = calls
296 .iter()
297 .all(|call| call.frame.bounds.is_rows() && !call.frame.bounds.is_unbounded());
298 let bounds = &calls.super_rows_frame_bounds;
299 match (bounds.n_preceding_rows(), bounds.n_following_rows()) {
300 (Some(n_preceding), Some(n_following))
301 if all_frames_bounded_rows && !args.order_key_indices.is_empty() =>
302 {
303 Some(StateCleaning {
304 watermark_col_idx: args.order_key_indices[0],
305 stale_rows_at_front: args.order_key_order_types[0].is_ascending(),
306 n_retain: n_preceding.saturating_add(n_following),
307 })
308 }
309 _ => {
310 tracing::warn!(
312 "state cleaning is enabled for over window with unbounded or non-`ROWS` frames, ignoring"
313 );
314 None
315 }
316 }
317 } else {
318 None
319 };
320
321 let watermark_forward_cols = {
322 let mut cols: HashSet<usize> = deduped_part_key_indices.iter().copied().collect();
323 if let Some(&first_order_key_idx) = args.order_key_indices.first()
324 && can_forward_watermark_on_order_key(
325 calls.iter().map(|call| &call.frame),
326 args.order_key_order_types[0],
327 )
328 {
329 cols.insert(first_order_key_idx);
330 }
331 cols
332 };
333
334 Self {
335 input: args.input,
336 inner: ExecutorInner {
337 actor_ctx: args.actor_ctx,
338 schema: args.schema,
339 calls,
340 deduped_part_key_indices,
341 order_key_indices: args.order_key_indices,
342 order_key_data_types,
343 order_key_order_types: args.order_key_order_types,
344 input_stream_key: input_info.stream_key,
345 state_key_to_table_sub_pk_proj,
346 state_table: args.state_table,
347 watermark_sequence: args.watermark_epoch,
348 chunk_size: args.chunk_size,
349 cache_policy,
350 state_cleaning,
351 watermark_forward_cols,
352 },
353 }
354 }
355
356 fn merge_changes_in_chunk<'a>(
363 this: &'_ ExecutorInner<S>,
364 chunk: &'a StreamChunk,
365 ) -> impl Iterator<Item = Record<RowRef<'a>>> {
366 let mut cb = ChangeBuffer::with_capacity(chunk.cardinality());
367 for record in chunk.records() {
368 cb.apply_record(record, |row| this.get_input_pk(row));
369 }
370 cb.into_records()
371 }
372
373 #[try_stream(ok = StreamChunk, error = StreamExecutorError)]
374 async fn apply_chunk<'a>(
375 this: &'a mut ExecutorInner<S>,
376 vars: &'a mut ExecutionVars<S>,
377 chunk: StreamChunk,
378 metrics: &'a OverWindowMetrics,
379 ) {
380 let mut deltas: BTreeMap<DefaultOrdered<OwnedRow>, (PartitionDelta, PartitionDelta)> =
385 BTreeMap::new();
386 let mut key_change_updated_pks = HashSet::new();
388
389 for record in Self::merge_changes_in_chunk(this, &chunk) {
391 match record {
392 Record::Insert { new_row } => {
393 let part_key = this.get_partition_key(new_row).into();
394 let (delta, _) = deltas.entry(part_key).or_default();
395 delta.insert(
396 this.row_to_cache_key(new_row)?,
397 Change::Insert(new_row.into_owned_row()),
398 );
399 }
400 Record::Delete { old_row } => {
401 let part_key = this.get_partition_key(old_row).into();
402 let (delta, _) = deltas.entry(part_key).or_default();
403 delta.insert(this.row_to_cache_key(old_row)?, Change::Delete);
404 }
405 Record::Update { old_row, new_row } => {
406 let old_part_key = this.get_partition_key(old_row).into();
407 let new_part_key = this.get_partition_key(new_row).into();
408 let old_state_key = this.row_to_cache_key(old_row)?;
409 let new_state_key = this.row_to_cache_key(new_row)?;
410 if old_part_key == new_part_key && old_state_key == new_state_key {
411 let (delta, no_effect_delta) = deltas.entry(old_part_key).or_default();
413 if old_row.project(&this.calls.all_arg_indices)
414 == new_row.project(&this.calls.all_arg_indices)
415 {
416 no_effect_delta
418 .insert(old_state_key, Change::Insert(new_row.into_owned_row()));
419 } else {
420 delta.insert(old_state_key, Change::Insert(new_row.into_owned_row()));
421 }
422 } else if old_part_key == new_part_key {
423 key_change_updated_pks.insert(this.get_input_pk(old_row));
426 let (delta, _) = deltas.entry(old_part_key).or_default();
427 delta.insert(old_state_key, Change::Delete);
428 delta.insert(new_state_key, Change::Insert(new_row.into_owned_row()));
429 } else {
430 let (old_part_delta, _) = deltas.entry(old_part_key).or_default();
435 old_part_delta.insert(old_state_key, Change::Delete);
436 let (new_part_delta, _) = deltas.entry(new_part_key).or_default();
437 new_part_delta
438 .insert(new_state_key, Change::Insert(new_row.into_owned_row()));
439 }
440 }
441 }
442 }
443
444 let mut key_change_update_buffer: BTreeMap<DefaultOrdered<OwnedRow>, Record<OwnedRow>> =
446 BTreeMap::new();
447 let mut chunk_builder = StreamChunkBuilder::new(this.chunk_size, this.schema.data_types());
448
449 for (part_key, (delta, no_effect_delta)) in deltas {
451 vars.stats.cache_lookup += 1;
452 if !vars.cached_partitions.contains(&part_key.0) {
453 vars.stats.cache_miss += 1;
454 vars.cached_partitions
455 .put(part_key.0.clone(), PartitionCache::new());
456 }
457 let mut cache = vars.cached_partitions.get_mut(&part_key).unwrap();
458
459 for (key, change) in no_effect_delta {
464 let new_row = change.into_insert().unwrap(); let (old_row, from_cache) = if let Some(old_row) = cache.inner().get(&key).cloned()
467 {
468 (old_row, true)
470 } else {
471 let table_pk = (&new_row).project(this.state_table.pk_indices());
473 if let Some(old_row) = this.state_table.get_row(table_pk).await? {
476 (old_row, false)
477 } else {
478 consistency_panic!(?part_key, ?key, ?new_row, "updating non-existing row");
479 continue;
480 }
481 };
482
483 let input_len = new_row.len();
485 let new_row = OwnedRow::new(
486 new_row
487 .into_iter()
488 .chain(old_row.as_inner().iter().skip(input_len).cloned()) .collect(),
490 );
491
492 let record = Record::Update {
494 old_row: &old_row,
495 new_row: &new_row,
496 };
497 if let Some(chunk) = chunk_builder.append_record(record.as_ref()) {
498 yield chunk;
499 }
500 this.state_table.write_record(record);
501 if from_cache {
502 cache.insert(key, new_row);
503 }
504 }
505
506 let mut partition = OverPartition::new(
507 &part_key,
508 &mut cache,
509 this.cache_policy,
510 &this.calls,
511 RowConverter {
512 state_key_to_table_sub_pk_proj: &this.state_key_to_table_sub_pk_proj,
513 order_key_indices: &this.order_key_indices,
514 order_key_data_types: &this.order_key_data_types,
515 order_key_order_types: &this.order_key_order_types,
516 input_stream_key_indices: &this.input_stream_key,
517 },
518 );
519
520 if delta.is_empty() {
521 continue;
522 }
523
524 if this.state_cleaning.is_some() {
525 vars.touched_partitions.insert(part_key.0.clone());
526 }
527
528 let (part_changes, accessed_range) =
530 partition.build_changes(&this.state_table, delta).await?;
531
532 for (key, record) in part_changes {
533 if !key_change_updated_pks.contains(&key.pk) {
535 if let Some(chunk) = chunk_builder.append_record(record.as_ref()) {
536 yield chunk;
537 }
538 } else {
539 let pk = key.pk.clone();
542 let record = record.clone();
543 if let Some(existed) = key_change_update_buffer.remove(&key.pk) {
544 match (existed, record) {
545 (Record::Insert { new_row }, Record::Delete { old_row })
546 | (Record::Delete { old_row }, Record::Insert { new_row }) => {
547 if let Some(chunk) =
549 chunk_builder.append_record(Record::Update { old_row, new_row })
550 {
551 yield chunk;
552 }
553 }
554 (existed, record) => {
555 consistency_panic!(
557 ?existed,
558 ?record,
559 "other cases should not exist",
560 );
561
562 key_change_update_buffer.insert(pk, record);
563 if let Some(chunk) = chunk_builder.append_record(existed) {
564 yield chunk;
565 }
566 }
567 }
568 } else {
569 key_change_update_buffer.insert(pk, record);
570 }
571 }
572
573 partition.write_record(&mut this.state_table, key, record);
575 }
576
577 if !key_change_update_buffer.is_empty() {
578 consistency_panic!(
579 ?key_change_update_buffer,
580 "key-change update buffer should be empty after processing"
581 );
582 }
585
586 let cache_len = partition.cache_real_len();
587 let stats = partition.summarize();
588 metrics
589 .over_window_range_cache_entry_count
590 .set(cache_len as i64);
591 metrics
592 .over_window_range_cache_lookup_count
593 .inc_by(stats.lookup_count);
594 metrics
595 .over_window_range_cache_left_miss_count
596 .inc_by(stats.left_miss_count);
597 metrics
598 .over_window_range_cache_right_miss_count
599 .inc_by(stats.right_miss_count);
600 metrics
601 .over_window_accessed_entry_count
602 .inc_by(stats.accessed_entry_count);
603 metrics
604 .over_window_compute_count
605 .inc_by(stats.compute_count);
606 metrics
607 .over_window_same_output_count
608 .inc_by(stats.same_output_count);
609
610 if !this.cache_policy.is_full()
612 && let Some(accessed_range) = accessed_range
613 {
614 match vars.recently_accessed_ranges.entry(part_key) {
615 btree_map::Entry::Vacant(vacant) => {
616 vacant.insert(accessed_range);
617 }
618 btree_map::Entry::Occupied(mut occupied) => {
619 let recently_accessed_range = occupied.get_mut();
620 let min_start = accessed_range
621 .start()
622 .min(recently_accessed_range.start())
623 .clone();
624 let max_end = accessed_range
625 .end()
626 .max(recently_accessed_range.end())
627 .clone();
628 *recently_accessed_range = min_start..=max_end;
629 }
630 }
631 }
632 }
633
634 if let Some(chunk) = chunk_builder.take() {
636 yield chunk;
637 }
638 }
639
640 async fn clean_state(
644 this: &mut ExecutorInner<S>,
645 vars: &mut ExecutionVars<S>,
646 ) -> StreamExecutorResult<usize> {
647 let touched = std::mem::take(&mut vars.touched_partitions);
648 let (Some(cleaning), Some(watermark)) = (&this.state_cleaning, &vars.cleaning_watermark)
649 else {
650 return Ok(0);
651 };
652 if touched.is_empty() {
653 return Ok(0);
654 }
655
656 let row_conv = RowConverter {
657 state_key_to_table_sub_pk_proj: &this.state_key_to_table_sub_pk_proj,
658 order_key_indices: &this.order_key_indices,
659 order_key_data_types: &this.order_key_data_types,
660 order_key_order_types: &this.order_key_order_types,
661 input_stream_key_indices: &this.input_stream_key,
662 };
663
664 let mut n_deleted = 0;
665 for part_key in touched {
666 let mut cache_guard = vars.cached_partitions.get_mut(&part_key);
667 let mut temp_cache = PartitionCache::new();
670 let cache = match cache_guard.as_deref_mut() {
671 Some(cache) => cache,
672 None => &mut temp_cache,
673 };
674 let mut partition =
675 OverPartition::new(&part_key, cache, this.cache_policy, &this.calls, row_conv);
676 let (n, has_more) = partition
677 .clean_stale_rows(&mut this.state_table, cleaning, watermark)
678 .await?;
679 n_deleted += n;
680 if has_more {
681 vars.touched_partitions.insert(part_key.clone());
683 }
684 }
685 Ok(n_deleted)
686 }
687
688 #[try_stream(ok = Message, error = StreamExecutorError)]
689 async fn executor_inner(self) {
690 let OverWindowExecutor {
691 input,
692 inner: mut this,
693 } = self;
694
695 let metrics_info = MetricsInfo::new(
696 this.actor_ctx.streaming_metrics.clone(),
697 this.state_table.table_id(),
698 this.actor_ctx.id,
699 "OverWindow",
700 );
701
702 let metrics = metrics_info.metrics.new_over_window_metrics(
703 this.state_table.table_id(),
704 this.actor_ctx.id,
705 this.actor_ctx.fragment_id,
706 );
707
708 let mut vars = ExecutionVars {
709 cached_partitions: ManagedLruCache::unbounded(
710 this.watermark_sequence.clone(),
711 metrics_info,
712 ),
713 recently_accessed_ranges: Default::default(),
714 cleaning_watermark: None,
715 touched_partitions: Default::default(),
716 stats: Default::default(),
717 _phantom: PhantomData::<S>,
718 };
719
720 let mut input = input.execute();
721 let barrier = expect_first_barrier(&mut input).await?;
722 let first_epoch = barrier.epoch;
723 yield Message::Barrier(barrier);
724 this.state_table.init_epoch(first_epoch).await?;
725
726 #[for_await]
727 for msg in input {
728 let msg = msg?;
729 match msg {
730 Message::Watermark(watermark) => {
731 if let Some(cleaning) = &this.state_cleaning
732 && watermark.col_idx == cleaning.watermark_col_idx
733 && vars
734 .cleaning_watermark
735 .as_ref()
736 .is_none_or(|old| old.default_cmp(&watermark.val).is_lt())
737 {
738 vars.cleaning_watermark = Some(watermark.val.clone());
740 }
741 if this.watermark_forward_cols.contains(&watermark.col_idx) {
742 yield Message::Watermark(watermark);
746 }
747 }
748 Message::Chunk(chunk) => {
749 #[for_await]
750 for chunk in Self::apply_chunk(&mut this, &mut vars, chunk, &metrics) {
751 yield Message::Chunk(chunk?);
752 }
753 this.state_table.try_flush().await?;
754
755 vars.cached_partitions.evict();
763 }
764 Message::Barrier(barrier) => {
765 let n_cleaned = Self::clean_state(&mut this, &mut vars).await?;
766 metrics
767 .over_window_state_cleaned_row_count
768 .inc_by(n_cleaned as u64);
769
770 let post_commit = this.state_table.commit(barrier.epoch).await?;
771
772 let update_vnode_bitmap = barrier.as_update_vnode_bitmap(this.actor_ctx.id);
773 yield Message::Barrier(barrier);
774
775 vars.cached_partitions.evict();
776
777 metrics
778 .over_window_cached_entry_count
779 .set(vars.cached_partitions.len() as _);
780 metrics
781 .over_window_cache_lookup_count
782 .inc_by(std::mem::take(&mut vars.stats.cache_lookup));
783 metrics
784 .over_window_cache_miss_count
785 .inc_by(std::mem::take(&mut vars.stats.cache_miss));
786
787 if let Some((_, cache_may_stale)) =
788 post_commit.post_yield_barrier(update_vnode_bitmap).await?
789 && cache_may_stale
790 {
791 vars.cached_partitions.clear();
792 vars.recently_accessed_ranges.clear();
793 vars.touched_partitions.clear();
794 }
795
796 if !this.cache_policy.is_full() {
797 for (part_key, recently_accessed_range) in
798 std::mem::take(&mut vars.recently_accessed_ranges)
799 {
800 if let Some(mut range_cache) =
801 vars.cached_partitions.get_mut(&part_key.0)
802 {
803 range_cache.shrink(
804 &part_key.0,
805 this.cache_policy,
806 recently_accessed_range,
807 );
808 }
809 }
810 }
811 }
812 }
813 }
814 }
815}
816
817#[derive(Debug, Clone, Copy)]
830pub(super) struct RowConverter<'a> {
831 state_key_to_table_sub_pk_proj: &'a [usize],
832 order_key_indices: &'a [usize],
833 order_key_data_types: &'a [DataType],
834 order_key_order_types: &'a [OrderType],
835 input_stream_key_indices: &'a [usize],
836}
837
838impl<'a> RowConverter<'a> {
839 pub(super) fn calc_state_key_to_table_sub_pk_proj(
843 partition_key_indices: &[usize],
844 order_key_indices: &[usize],
845 input_stream_key_indices: &'a [usize],
846 ) -> Vec<usize> {
847 let mut projection =
849 Vec::with_capacity(order_key_indices.len() + input_stream_key_indices.len());
850 let mut col_dedup: HashSet<usize> = partition_key_indices.iter().copied().collect();
851 for (proj_idx, key_idx) in order_key_indices
852 .iter()
853 .chain(input_stream_key_indices.iter())
854 .enumerate()
855 {
856 if col_dedup.insert(*key_idx) {
857 projection.push(proj_idx);
858 }
859 }
860 projection.shrink_to_fit();
861 projection
862 }
863
864 pub(super) fn state_key_to_table_sub_pk(
866 &self,
867 key: &StateKey,
868 ) -> StreamExecutorResult<OwnedRow> {
869 Ok(memcmp_encoding::decode_row(
870 &key.order_key,
871 self.order_key_data_types,
872 self.order_key_order_types,
873 )?
874 .chain(key.pk.as_inner())
875 .project(self.state_key_to_table_sub_pk_proj)
876 .into_owned_row())
877 }
878
879 pub(super) fn row_to_state_key(
881 &self,
882 full_row: impl Row + Copy,
883 ) -> StreamExecutorResult<StateKey> {
884 Ok(StateKey {
885 order_key: memcmp_encoding::encode_row(
886 full_row.project(self.order_key_indices),
887 self.order_key_order_types,
888 )?,
889 pk: full_row
890 .project(self.input_stream_key_indices)
891 .into_owned_row()
892 .into(),
893 })
894 }
895}
896
897#[cfg(test)]
898mod tests {
899 use std::ops::Bound;
900
901 use futures::TryStreamExt;
902 use risingwave_common::catalog::{ColumnDesc, ColumnId, TableId};
903 use risingwave_common::util::epoch::{EpochPair, test_epoch};
904 use risingwave_storage::memory::MemoryStateStore;
905 use risingwave_storage::store::PrefetchOptions;
906
907 use super::*;
908 use crate::common::table::test_utils::gen_pbtable;
909
910 #[tokio::test]
911 async fn test_state_cleaning_large_retention() {
912 for order_type in [OrderType::ascending(), OrderType::descending()] {
913 for cached in [true, false] {
914 for n_retain in [usize::MAX - 65_534, usize::MAX - 1] {
918 let mut table = StateTable::from_table_catalog(
919 &gen_pbtable(
920 TableId::new(1),
921 vec![ColumnDesc::unnamed(ColumnId::new(0), DataType::Int64)],
922 vec![order_type],
923 vec![0],
924 0,
925 ),
926 MemoryStateStore::new(),
927 None,
928 )
929 .await;
930 table
931 .init_epoch(EpochPair::new_test_epoch(test_epoch(1)))
932 .await
933 .unwrap();
934 let rows = [10i64, 20, 30].map(|value| OwnedRow::new(vec![Some(value.into())]));
935 let row_conv = RowConverter {
936 state_key_to_table_sub_pk_proj: &[0],
937 order_key_indices: &[0],
938 order_key_data_types: &[DataType::Int64],
939 order_key_order_types: &[order_type],
940 input_stream_key_indices: &[0],
941 };
942 let mut cache = if cached {
943 PartitionCache::new_without_sentinels()
944 } else {
945 PartitionCache::new()
946 };
947 for row in &rows {
948 table.insert(row.clone());
949 if cached {
950 cache.insert(
951 CacheKey::from(row_conv.row_to_state_key(row).unwrap()),
952 row.clone(),
953 );
954 }
955 }
956 table
957 .commit_for_test(EpochPair::new_test_epoch(test_epoch(2)))
958 .await
959 .unwrap();
960
961 let partition_key = OwnedRow::empty();
962 let calls = Calls::new(vec![]);
963 let mut partition = OverPartition::new(
964 &partition_key,
965 &mut cache,
966 CachePolicy::Full,
967 &calls,
968 row_conv,
969 );
970 let cleaning = StateCleaning {
971 watermark_col_idx: 0,
972 stale_rows_at_front: order_type.is_ascending(),
973 n_retain,
974 };
975 assert_eq!(
976 partition
977 .clean_stale_rows(&mut table, &cleaning, &100i64.into())
978 .await
979 .unwrap(),
980 (0, false),
981 "order={order_type:?}, cached={cached}, n_retain={n_retain}",
982 );
983 table
984 .commit_for_test(EpochPair::new_test_epoch(test_epoch(3)))
985 .await
986 .unwrap();
987 let range: (Bound<OwnedRow>, Bound<OwnedRow>) =
988 (Bound::Unbounded, Bound::Unbounded);
989 let remaining: Vec<OwnedRow> = table
990 .iter_with_prefix(&partition_key, &range, PrefetchOptions::default())
991 .await
992 .unwrap()
993 .try_collect()
994 .await
995 .unwrap();
996 let mut expected = rows.to_vec();
997 if !order_type.is_ascending() {
998 expected.reverse();
999 }
1000 assert_eq!(remaining, expected);
1001 assert_eq!(cache.normal_len(), if cached { rows.len() } else { 0 });
1002 }
1003 }
1004 }
1005 }
1006}