Skip to main content

risingwave_stream/executor/over_window/
general.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, 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
41/// [`OverWindowExecutor`] consumes retractable input stream and produces window function outputs.
42/// One [`OverWindowExecutor`] can handle one combination of partition key and order key.
43///
44/// - State table schema = output schema, state table pk = `partition key | order key | input pk`.
45/// - Output schema = input schema + window function results.
46/// - Watermarks on partition key columns are always forwarded, since a row can only affect rows in
47///   the same partition. Watermarks on the first order key column are forwarded only if the window
48///   frames guarantee that a row can only affect rows not below itself in that column (see
49///   [`can_forward_watermark_on_order_key`]). All other watermarks are dropped.
50/// - When [`StateCleaning`] is enabled, stale rows below the watermark of the first order key
51///   column are deleted from recently touched partitions at barriers.
52pub 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    /// The maximum size of the chunk produced by executor at a time.
73    chunk_size: usize,
74    cache_policy: CachePolicy,
75    /// Watermark-driven state cleaning strategy, `None` if disabled.
76    state_cleaning: Option<StateCleaning>,
77    /// Indices of input columns on which watermarks are forwarded to downstream: the partition key
78    /// columns, and the first order key column if allowed by the window frames.
79    watermark_forward_cols: HashSet<usize>,
80}
81
82struct ExecutionVars<S: StateStore> {
83    /// partition key => partition range cache.
84    cached_partitions: ManagedLruCache<OwnedRow, PartitionCache>,
85    /// partition key => recently accessed range.
86    recently_accessed_ranges: BTreeMap<DefaultOrdered<OwnedRow>, RangeInclusive<StateKey>>,
87    /// The latest watermark received on the watermark column for state cleaning.
88    cleaning_watermark: Option<ScalarImpl>,
89    /// Partitions touched since the last barrier, which need state cleaning at the next barrier.
90    touched_partitions: HashSet<OwnedRow>,
91    stats: ExecutionStats,
92    _phantom: PhantomData<S>,
93}
94
95/// Watermark-driven state cleaning strategy of [`OverWindowExecutor`].
96///
97/// This is only enabled when the optimizer decides it's safe, i.e., when the input is append-only,
98/// all window frames are bounded `ROWS` frames, and the first order key column is a watermark
99/// column with NULLs ordered as the largest values.
100///
101/// Under these conditions, a row can only affect (and be affected by) a bounded number of
102/// neighboring rows in the partition. Once a watermark `wm` is received, no row with order key
103/// `< wm` will ever arrive, so rows with order key `< wm` (*stale rows*) can never get new
104/// neighbors on the "smaller" side. Therefore, among the stale rows, only the `n_retain` ones that
105/// are closest to the watermark boundary can still be involved in future computation (as frame
106/// members of, or as rows affected by, rows arriving in the future), and all the others can be
107/// safely deleted from the state table.
108///
109/// The cleaning happens at barriers for partitions touched since the last barrier. This keeps the
110/// work proportional to recently active partitions and avoids keeping per-partition cleanup
111/// metadata in memory. Active partitions retain only the required rows plus rows not yet behind
112/// the watermark. A partition that stops receiving rows is cleaned lazily when it's touched again.
113#[derive(Debug)]
114pub(super) struct StateCleaning {
115    /// Index of the watermark column (the first order key column) in the input schema.
116    pub watermark_col_idx: usize,
117    /// Whether stale rows are at the front (`ASC` order key) or at the back (`DESC` order key)
118    /// of the partition.
119    pub stale_rows_at_front: bool,
120    /// Number of stale rows to retain in each partition, which is the number of preceding rows
121    /// plus the number of following rows of the union of all `ROWS` frames.
122    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    /// Get deduplicated partition key from a full row, which happened to be the prefix of table PK.
139    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    /// `full_row` can be an input row or state table row.
150    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    /// Whether to enable watermark-driven state cleaning. See [`StateCleaning`].
183    pub enable_state_cleaning: bool,
184}
185
186/// Information about the window function calls.
187/// Contains the original calls and many other information that can be derived from the calls to avoid
188/// repeated calculation.
189pub(super) struct Calls {
190    calls: Vec<WindowFuncCall>,
191
192    /// The `ROWS` frame that is the union of all `ROWS` frames.
193    pub(super) super_rows_frame_bounds: RowsFrameBounds,
194    /// All `RANGE` frames.
195    pub(super) range_frames: Vec<RangeFrameBounds>,
196    pub(super) start_is_unbounded: bool,
197    pub(super) end_is_unbounded: bool,
198    /// Deduplicated indices of all arguments of all calls.
199    pub(super) all_arg_indices: Vec<usize>,
200
201    // TODO(rc): The following flags are used to optimize for `row_number`, `rank` and `dense_rank`.
202    // We should try our best to remove these flags while maintaining the performance in the future.
203    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            // For unbounded frames, we finally need all entries of the partition in the cache,
267            // so for simplicity we just use full cache policy for these cases.
268            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                    // The optimizer should never enable state cleaning in this case.
311                    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    /// Merge changes by input pk in the given chunk, return a change iterator which guarantees that
357    /// each pk only appears once. This method also validates the consistency of the input
358    /// chunk.
359    ///
360    /// TODO(rc): We may want to optimize this by handling changes on the same pk during generating
361    /// partition [`Change`]s.
362    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        // (deduped) partition key => (
381        //   significant changes happened in the partition,
382        //   no-effect changes happened in the partition,
383        // )
384        let mut deltas: BTreeMap<DefaultOrdered<OwnedRow>, (PartitionDelta, PartitionDelta)> =
385            BTreeMap::new();
386        // input pk of update records of which the order key is changed.
387        let mut key_change_updated_pks = HashSet::new();
388
389        // Collect changes for each partition.
390        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                        // not a key-change update
412                        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                            // partition key, order key and arguments are all the same
417                            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                        // order-change update, split into delete + insert, will be merged after
424                        // building changes
425                        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                        // partition-change update, split into delete + insert
431                        // NOTE(rc): Since we append partition key to logical pk, we can't merge the
432                        // delete + insert back to update later.
433                        // TODO: IMO this behavior is problematic. Deep discussion is needed.
434                        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        // `input pk` => `Record`
445        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        // Build final changes partition by partition.
450        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            // First, handle `Update`s that don't affect window function outputs.
460            // Be careful that changes in `delta` may (though we believe unlikely) affect the
461            // window function outputs of rows in `no_effect_delta`, so before handling `delta`
462            // we need to write all changes to state table, range cache and chunk builder.
463            for (key, change) in no_effect_delta {
464                let new_row = change.into_insert().unwrap(); // new row of an `Update`
465
466                let (old_row, from_cache) = if let Some(old_row) = cache.inner().get(&key).cloned()
467                {
468                    // Got old row from range cache.
469                    (old_row, true)
470                } else {
471                    // Retrieve old row from state table.
472                    let table_pk = (&new_row).project(this.state_table.pk_indices());
473                    // The accesses to the state table is ordered by table PK, so ideally we
474                    // can leverage the block cache under the hood.
475                    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                // concatenate old outputs
484                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()) // chain old outputs
489                        .collect(),
490                );
491
492                // apply & emit the change
493                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            // Build changes for current partition.
529            let (part_changes, accessed_range) =
530                partition.build_changes(&this.state_table, delta).await?;
531
532            for (key, record) in part_changes {
533                // Build chunk and yield if needed.
534                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                    // For key-change updates, we should wait for both `Delete` and `Insert` changes
540                    // and merge them together.
541                    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                                // merge `Delete` and `Insert` into `Update`
548                                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                                // when stream is inconsistent, there may be an `Update` of which the old pk does not actually exist
556                                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                // Apply the change record.
574                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                // if in non-strict mode, we can reach here, but we don't know the `StateKey`,
583                // so just ignore the buffer.
584            }
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            // Update recently accessed range for later shrinking cache.
611            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        // Yield remaining changes to downstream.
635        if let Some(chunk) = chunk_builder.take() {
636            yield chunk;
637        }
638    }
639
640    /// Clean up stale rows of the partitions touched since the last barrier, according to the
641    /// latest watermark received.
642    /// Returns the number of rows deleted. See [`StateCleaning`].
643    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            // If the partition is not cached (e.g. evicted), use a temporary cache with only
668            // sentinels, so that `OverPartition` scans the state table for stale rows.
669            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                // continue to clean this partition at the next barrier
682                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                        // Only used for state cleaning at the next barrier.
739                        vars.cleaning_watermark = Some(watermark.val.clone());
740                    }
741                    if this.watermark_forward_cols.contains(&watermark.col_idx) {
742                        // All changes caused by the rows received so far have already been
743                        // emitted, and rows arriving in the future can only affect rows that are
744                        // not below the watermark in this column, so it's safe to forward it now.
745                        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                    // Also apply the LRU watermark at chunk boundaries, so that cold
756                    // partitions can be released without waiting for the next barrier,
757                    // which can be a long time away with large barrier intervals. This
758                    // is safe because the range cache is write-through: at this point
759                    // all changes have been applied to both the state table and the
760                    // cache, so an evicted partition can be reloaded from the state
761                    // table with identical content.
762                    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/// A converter that helps convert [`StateKey`] to state table sub-PK and convert executor input/output
818/// row to [`StateKey`].
819///
820/// ## Notes
821///
822/// - [`StateKey`]: Over window range cache key type, containing order key and input pk.
823/// - State table sub-PK: State table PK = PK prefix (partition key) + sub-PK (order key + input pk).
824/// - Input/output row: Input schema is the prefix of output schema.
825///
826/// You can see that the content of [`StateKey`] is very similar to state table sub-PK. There's only
827/// one difference: the state table PK and sub-PK don't have duplicated columns, while in [`StateKey`],
828/// `order_key` and (input)`pk` may contain duplicated columns.
829#[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    /// Calculate the indices needed for projection from [`StateKey`] to state table sub-PK (used to do
840    /// prefixed table scanning). Ideally this function should be called only once by each executor instance.
841    /// The projection indices vec is the *selected column indices* in [`StateKey`].`order_key.chain(input_pk)`.
842    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        // This process is corresponding to `StreamOverWindow::infer_state_table`.
848        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    /// Convert [`StateKey`] to sub-PK (table PK without partition key) as [`OwnedRow`].
865    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    /// Convert full input/output row to [`StateKey`].
880    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                // Both are valid sums of two Int64 ROWS offsets on 64-bit targets.
915                // The first makes the old collection limit wrap to 1, so even three
916                // rows catch the erroneous deletion with overflow checks disabled.
917                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}