risingwave_common/array/
stream_chunk_iter.rs1use std::ops::Range;
16
17use super::RowRef;
18use super::data_chunk_iter::DataChunkRefIter;
19use super::stream_record::Record;
20use crate::array::{Op, StreamChunk};
21use crate::row::RowExt;
22
23impl StreamChunk {
24 pub fn records(&self) -> StreamChunkRefIter<'_> {
26 StreamChunkRefIter {
27 chunk: self,
28 inner: self.data_chunk().rows(),
29 }
30 }
31
32 pub fn rows(&self) -> impl Iterator<Item = (Op, RowRef<'_>)> {
36 self.rows_in(0..self.capacity())
37 }
38
39 pub fn rows_in(&self, range: Range<usize>) -> impl Iterator<Item = (Op, RowRef<'_>)> {
41 self.data_chunk().rows_in(range).map(|row| {
42 (
43 unsafe { *self.ops().get_unchecked(row.index()) },
45 row,
46 )
47 })
48 }
49
50 pub fn row_at(&self, pos: usize) -> (Op, RowRef<'_>, bool) {
52 let op = self.ops()[pos];
53 let (row, visible) = self.data_chunk().row_at(pos);
54 (op, row, visible)
55 }
56
57 pub fn rows_with_holes(&self) -> impl ExactSizeIterator<Item = Option<(Op, RowRef<'_>)>> {
59 self.data_chunk().rows_with_holes().map(|row| {
60 row.map(|row| {
61 (
62 unsafe { *self.ops().get_unchecked(row.index()) },
64 row,
65 )
66 })
67 })
68 }
69}
70
71pub struct StreamChunkRefIter<'a> {
72 chunk: &'a StreamChunk,
73
74 inner: DataChunkRefIter<'a>,
75}
76
77impl<'a> Iterator for StreamChunkRefIter<'a> {
78 type Item = Record<RowRef<'a>>;
79
80 fn next(&mut self) -> Option<Self::Item> {
81 let row = self.inner.next()?;
82 let op = unsafe { self.chunk.ops().get_unchecked(row.index()) };
84
85 match op {
86 Op::Insert => Some(Record::Insert { new_row: row }),
87 Op::Delete => Some(Record::Delete { old_row: row }),
88 Op::UpdateDelete => {
89 let next_row = self
90 .inner
91 .next()
92 .unwrap_or_else(|| panic!("expect a row after U-\nU- row: {}", row.display()));
93 let op = unsafe { *self.chunk.ops().get_unchecked(next_row.index()) };
95 debug_assert_eq!(
96 op,
97 Op::UpdateInsert,
98 "expect a U+ after U-\nU- row: {}\nrow after U-: {}",
99 row.display(),
100 next_row.display()
101 );
102
103 Some(Record::Update {
104 old_row: row,
105 new_row: next_row,
106 })
107 }
108 Op::UpdateInsert => panic!("expect a U- before U+\nU+ row: {}", row.display()),
109 }
110 }
111
112 fn size_hint(&self) -> (usize, Option<usize>) {
113 let (lower, upper) = self.inner.size_hint();
114 (lower / 2, upper)
115 }
116}
117
118#[cfg(test)]
119mod tests {
120 extern crate test;
121 use itertools::Itertools;
122 use test::Bencher;
123
124 use crate::array::stream_record::Record;
125 use crate::row::Row;
126 use crate::test_utils::test_stream_chunk::{
127 BigStreamChunk, TestStreamChunk, WhatEverStreamChunk,
128 };
129
130 #[test]
131 fn test_chunk_rows() {
132 let test = WhatEverStreamChunk;
133 let chunk = test.stream_chunk();
134 let mut rows = chunk.rows().map(|(op, row)| (op, row.into_owned_row()));
135 assert_eq!(Some(test.row_with_op_at(0)), rows.next());
136 assert_eq!(Some(test.row_with_op_at(1)), rows.next());
137 assert_eq!(Some(test.row_with_op_at(2)), rows.next());
138 assert_eq!(Some(test.row_with_op_at(3)), rows.next());
139 }
140
141 #[test]
142 fn test_chunk_records() {
143 let test = WhatEverStreamChunk;
144 let chunk = test.stream_chunk();
145 let mut rows = chunk
146 .records()
147 .flat_map(Record::into_rows)
148 .map(|(op, row)| (op, row.into_owned_row()));
149 assert_eq!(Some(test.row_with_op_at(0)), rows.next());
150 assert_eq!(Some(test.row_with_op_at(1)), rows.next());
151 assert_eq!(Some(test.row_with_op_at(2)), rows.next());
152 assert_eq!(Some(test.row_with_op_at(3)), rows.next());
153 }
154
155 #[bench]
156 fn bench_rows_iterator_from_records(b: &mut Bencher) {
157 let chunk = BigStreamChunk::new(10000).stream_chunk();
158 b.iter(|| {
159 for (_op, row) in chunk.records().flat_map(Record::into_rows) {
160 test::black_box(row.iter().count());
161 }
162 })
163 }
164
165 #[bench]
166 fn bench_rows_iterator(b: &mut Bencher) {
167 let chunk = BigStreamChunk::new(10000).stream_chunk();
168 b.iter(|| {
169 for (_op, row) in chunk.rows() {
170 test::black_box(row.iter().count());
171 }
172 })
173 }
174
175 #[bench]
176 fn bench_rows_iterator_vec_of_datum_refs(b: &mut Bencher) {
177 let chunk = BigStreamChunk::new(10000).stream_chunk();
178 b.iter(|| {
179 for (_op, row) in chunk.rows() {
180 let row = row.iter().collect_vec();
182 test::black_box(row);
183 }
184 })
185 }
186
187 #[bench]
188 fn bench_record_iterator(b: &mut Bencher) {
189 let chunk = BigStreamChunk::new(10000).stream_chunk();
190 b.iter(|| {
191 for record in chunk.records() {
192 match record {
193 Record::Insert { new_row } => test::black_box(new_row.iter().count()),
194 Record::Delete { old_row } => test::black_box(old_row.iter().count()),
195 Record::Update { old_row, new_row } => {
196 test::black_box(old_row.iter().count());
197 test::black_box(new_row.iter().count())
198 }
199 };
200 }
201 })
202 }
203}