Skip to main content

risingwave_frontend/optimizer/plan_node/
stream_topn.rs

1// Copyright 2022 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::assert_matches;
16
17use pretty_xmlish::XmlNode;
18use risingwave_common::util::sort_util::topn_watermark_forwardable_order_key;
19use risingwave_pb::stream_plan::stream_node::PbNodeBody;
20
21use super::generic::{DistillUnit, TopNLimit};
22use super::stream::prelude::*;
23use super::utils::{Distill, plan_node_name, watermark_pretty};
24use super::{
25    ExprRewritable, PlanBase, PlanTreeNodeUnary, StreamNode, StreamPlanRef as PlanRef, generic,
26};
27use crate::optimizer::plan_node::expr_visitable::ExprVisitable;
28use crate::optimizer::property::{Distribution, MonotonicityMap, Order, WatermarkColumns};
29use crate::stream_fragmenter::BuildFragmentGraphState;
30
31/// `StreamTopN` implements [`super::LogicalTopN`] to find the top N elements with a heap
32#[derive(Debug, Clone, PartialEq, Eq, Hash)]
33pub struct StreamTopN {
34    pub base: PlanBase<Stream>,
35    core: generic::TopN<PlanRef>,
36}
37
38impl StreamTopN {
39    pub fn new(core: generic::TopN<PlanRef>) -> Result<Self> {
40        assert!(core.group_key.is_empty());
41        assert!(core.limit_attr.limit() > 0);
42        let input = &core.input;
43        assert_matches!(input.distribution(), Distribution::Single);
44        reject_upsert_input!(input);
45        // The executor only forwards watermarks on the first `ORDER BY` column ordered
46        // `ASC NULLS LAST`. Keep the optimizer in sync so EOWC won't be enabled on unsupported
47        // plans. See `topn_watermark_forwardable_order_key` for the reasoning.
48        let watermark_columns =
49            match topn_watermark_forwardable_order_key(&core.order.column_orders) {
50                Some(col_idx) => input.watermark_columns().retain_clone(&[col_idx]),
51                None => WatermarkColumns::new(),
52            };
53
54        let base = PlanBase::new_stream_with_core(
55            &core,
56            Distribution::Single,
57            StreamKind::Retract,
58            false,
59            watermark_columns,
60            MonotonicityMap::new(),
61        );
62
63        Ok(StreamTopN { base, core })
64    }
65
66    pub fn limit_attr(&self) -> TopNLimit {
67        self.core.limit_attr
68    }
69
70    pub fn offset(&self) -> u64 {
71        self.core.offset
72    }
73
74    pub fn topn_order(&self) -> &Order {
75        &self.core.order
76    }
77}
78
79impl Distill for StreamTopN {
80    fn distill<'a>(&self) -> XmlNode<'a> {
81        let name = plan_node_name!("StreamTopN",
82            { "append_only", self.input().append_only() },
83        );
84        let mut node = self.core.distill_with_name(name);
85        if let Some(ow) = watermark_pretty(self.base.watermark_columns(), self.schema()) {
86            node.fields.push(("output_watermarks".into(), ow));
87        }
88        node
89    }
90}
91
92impl PlanTreeNodeUnary<Stream> for StreamTopN {
93    fn input(&self) -> PlanRef {
94        self.core.input.clone()
95    }
96
97    fn clone_with_input(&self, input: PlanRef) -> Self {
98        let mut core = self.core.clone();
99        core.input = input;
100        Self::new(core).unwrap()
101    }
102}
103
104impl_plan_tree_node_for_unary! { Stream, StreamTopN }
105
106impl StreamNode for StreamTopN {
107    fn to_stream_prost_body(&self, state: &mut BuildFragmentGraphState) -> PbNodeBody {
108        use risingwave_pb::stream_plan::*;
109
110        let input = self.input();
111        let topn_node = TopNNode {
112            limit: self.limit_attr().limit(),
113            offset: self.offset(),
114            with_ties: self.limit_attr().with_ties(),
115            table: Some(
116                self.core
117                    .infer_internal_table_catalog(
118                        input.schema(),
119                        input.ctx(),
120                        input.expect_stream_key(),
121                        None,
122                    )
123                    .with_id(state.gen_table_id_wrapped())
124                    .to_internal_table_prost(),
125            ),
126            order_by: self.topn_order().to_protobuf(),
127        };
128        if self.input().append_only() {
129            PbNodeBody::AppendOnlyTopN(Box::new(topn_node))
130        } else {
131            PbNodeBody::TopN(Box::new(topn_node))
132        }
133    }
134}
135
136impl ExprRewritable<Stream> for StreamTopN {}
137
138impl ExprVisitable for StreamTopN {}