Skip to main content

risingwave_common/array/
stream_chunk_iter.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::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    /// Return an iterator on stream records of this stream chunk.
25    pub fn records(&self) -> StreamChunkRefIter<'_> {
26        StreamChunkRefIter {
27            chunk: self,
28            inner: self.data_chunk().rows(),
29        }
30    }
31
32    /// Return an iterator on rows of this stream chunk.
33    ///
34    /// Should consider using [`StreamChunk::records`] if possible.
35    pub fn rows(&self) -> impl Iterator<Item = (Op, RowRef<'_>)> {
36        self.rows_in(0..self.capacity())
37    }
38
39    /// Return an iterator on rows of this stream chunk in a range.
40    pub fn rows_in(&self, range: Range<usize>) -> impl Iterator<Item = (Op, RowRef<'_>)> {
41        self.data_chunk().rows_in(range).map(|row| {
42            (
43                // SAFETY: index is checked since we are in the iterator.
44                unsafe { *self.ops().get_unchecked(row.index()) },
45                row,
46            )
47        })
48    }
49
50    /// Random access a row at `pos`. Return the op, data and whether the row is visible.
51    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    /// Iterates over every physical row position, yielding `None` for invisible rows.
58    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                    // SAFETY: index is checked since we are in the iterator.
63                    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        // SAFETY: index is checked since `row` is `Some`.
83        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                // SAFETY: index is checked since `insert_row` is `Some`.
94                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                // Mimic the old `RowRef(Vec<DatumRef>)`
181                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}