risingwave_frontend/optimizer/plan_node/
stream_topn.rs1use 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#[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 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 {}