Skip to main content

risingwave_common/array/
stream_chunk.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::fmt::Display;
16use std::marker::PhantomData;
17use std::mem::size_of;
18use std::ops::{Deref, DerefMut};
19use std::sync::Arc;
20use std::{fmt, mem};
21
22use either::Either;
23use enum_as_inner::EnumAsInner;
24use itertools::Itertools;
25use rand::prelude::SmallRng;
26use rand::{Rng, SeedableRng};
27use risingwave_common_estimate_size::EstimateSize;
28use risingwave_pb::data::{PbOp, PbStreamChunk};
29
30use super::stream_chunk_builder::StreamChunkBuilder;
31use super::{ArrayImpl, ArrayRef, ArrayResult, DataChunkTestExt, RowRef};
32use crate::array::DataChunk;
33use crate::bitmap::{Bitmap, BitmapBuilder};
34use crate::catalog::Schema;
35use crate::field_generator::VarcharProperty;
36use crate::row::Row;
37use crate::types::{DataType, DefaultOrdered, ToText};
38
39/// `Op` represents three operations in `StreamChunk`.
40///
41/// `UpdateDelete` and `UpdateInsert` are semantically equivalent to `Delete` and `Insert`
42/// but always appear in pairs to represent an update operation. It's guaranteed that
43/// they are adjacent to each other in the same `StreamChunk`, and the stream key of the two
44/// rows are the same.
45#[derive(Clone, Copy, Debug, PartialOrd, Ord, PartialEq, Eq, Hash, EnumAsInner)]
46pub enum Op {
47    Insert,
48    Delete,
49    UpdateDelete,
50    UpdateInsert,
51}
52
53impl Op {
54    pub fn to_protobuf(self) -> PbOp {
55        match self {
56            Op::Insert => PbOp::Insert,
57            Op::Delete => PbOp::Delete,
58            Op::UpdateInsert => PbOp::UpdateInsert,
59            Op::UpdateDelete => PbOp::UpdateDelete,
60        }
61    }
62
63    pub fn from_protobuf(prost: &i32) -> ArrayResult<Op> {
64        let op = match PbOp::try_from(*prost) {
65            Ok(PbOp::Insert) => Op::Insert,
66            Ok(PbOp::Delete) => Op::Delete,
67            Ok(PbOp::UpdateInsert) => Op::UpdateInsert,
68            Ok(PbOp::UpdateDelete) => Op::UpdateDelete,
69            Ok(PbOp::Unspecified) => unreachable!(),
70            Err(_) => bail!("No such op type"),
71        };
72        Ok(op)
73    }
74
75    /// convert `UpdateDelete` to `Delete` and `UpdateInsert` to Insert
76    pub fn normalize_update(self) -> Op {
77        match self {
78            Op::Insert => Op::Insert,
79            Op::Delete => Op::Delete,
80            Op::UpdateDelete => Op::Delete,
81            Op::UpdateInsert => Op::Insert,
82        }
83    }
84
85    pub fn to_i16(self) -> i16 {
86        match self {
87            Op::Insert => 1,
88            Op::Delete => 2,
89            Op::UpdateInsert => 3,
90            Op::UpdateDelete => 4,
91        }
92    }
93
94    pub fn to_varchar(self) -> String {
95        match self {
96            Op::Insert => "Insert",
97            Op::Delete => "Delete",
98            Op::UpdateInsert => "UpdateInsert",
99            Op::UpdateDelete => "UpdateDelete",
100        }
101        .to_owned()
102    }
103}
104
105/// `StreamChunk` is used to pass data over the streaming pathway.
106#[derive(Clone, PartialEq)]
107pub struct StreamChunk {
108    // TODO: Optimize using bitmap
109    ops: Arc<[Op]>,
110    data: DataChunk,
111}
112
113impl Default for StreamChunk {
114    /// Create a 0-row-0-col `StreamChunk`. Only used in some existing tests.
115    /// This is NOT the same as an **empty** chunk, which has 0 rows but with
116    /// columns aligned with executor schema.
117    fn default() -> Self {
118        Self {
119            ops: Arc::new([]),
120            data: DataChunk::new(vec![], 0),
121        }
122    }
123}
124
125impl StreamChunk {
126    /// Create a new `StreamChunk` with given ops and columns.
127    pub fn new(ops: impl Into<Arc<[Op]>>, columns: Vec<ArrayRef>) -> Self {
128        let ops = ops.into();
129        let visibility = Bitmap::ones(ops.len());
130        Self::with_visibility(ops, columns, visibility)
131    }
132
133    /// Create a new `StreamChunk` with given ops, columns and visibility.
134    pub fn with_visibility(
135        ops: impl Into<Arc<[Op]>>,
136        columns: Vec<ArrayRef>,
137        visibility: Bitmap,
138    ) -> Self {
139        let ops = ops.into();
140        for col in &columns {
141            assert_eq!(col.len(), ops.len());
142        }
143        let data = DataChunk::new(columns, visibility);
144        StreamChunk { ops, data }
145    }
146
147    /// Build a `StreamChunk` from rows.
148    ///
149    /// Panics if the `rows` is empty.
150    ///
151    /// Should prefer using [`StreamChunkBuilder`] instead to avoid unnecessary
152    /// allocation of rows.
153    pub fn from_rows(rows: &[(Op, impl Row)], data_types: &[DataType]) -> Self {
154        let mut builder = StreamChunkBuilder::unlimited(data_types.to_vec(), Some(rows.len()));
155
156        for (op, row) in rows {
157            let none = builder.append_row(*op, row);
158            debug_assert!(none.is_none());
159        }
160
161        builder.take().expect("chunk should not be empty")
162    }
163
164    pub fn empty(data_types: &[DataType]) -> Self {
165        StreamChunkBuilder::build_empty(data_types.to_vec())
166    }
167
168    /// Get the reference of the underlying data chunk.
169    pub fn data_chunk(&self) -> &DataChunk {
170        &self.data
171    }
172
173    /// Removes the invisible rows based on `visibility`. Returns a new compacted chunk
174    /// with all rows visible.
175    ///
176    /// This does not change the visible content of the chunk. Not to be confused with
177    /// `StreamChunkCompactor`, which removes unnecessary changes based on the key.
178    ///
179    /// See [`DataChunk::compact_vis`] for more details.
180    pub fn compact_vis(self) -> Self {
181        if self.is_vis_compacted() {
182            return self;
183        }
184
185        let (ops, columns, visibility) = self.into_inner();
186
187        let cardinality = visibility
188            .iter()
189            .fold(0, |vis_cnt, vis| vis_cnt + vis as usize);
190        let columns: Vec<_> = columns
191            .into_iter()
192            .map(|col| col.compact_vis(&visibility, cardinality).into())
193            .collect();
194        let mut new_ops = Vec::with_capacity(cardinality);
195        for idx in visibility.iter_ones() {
196            new_ops.push(ops[idx]);
197        }
198        StreamChunk::new(new_ops, columns)
199    }
200
201    /// Split the `StreamChunk` into multiple chunks with the given size at most.
202    ///
203    /// When the total cardinality of all the chunks is not evenly divided by the `size`,
204    /// the last new chunk will be the remainder.
205    ///
206    /// For consecutive `UpdateDelete` and `UpdateInsert`, they will be kept in one chunk.
207    /// As a result, some chunks may have `size + 1` rows.
208    pub fn split(&self, size: usize) -> Vec<Self> {
209        let mut builder = StreamChunkBuilder::new(size, self.data_types());
210        let mut outputs = Vec::new();
211
212        // TODO: directly append the chunk.
213        for (op, row) in self.rows() {
214            if let Some(chunk) = builder.append_row(op, row) {
215                outputs.push(chunk);
216            }
217        }
218        if let Some(output) = builder.take() {
219            outputs.push(output);
220        }
221
222        outputs
223    }
224
225    pub fn into_parts(self) -> (DataChunk, Arc<[Op]>) {
226        (self.data, self.ops)
227    }
228
229    pub fn from_parts(ops: impl Into<Arc<[Op]>>, data_chunk: DataChunk) -> Self {
230        let (columns, vis) = data_chunk.into_parts();
231        Self::with_visibility(ops, columns, vis)
232    }
233
234    pub fn into_inner(self) -> (Arc<[Op]>, Vec<ArrayRef>, Bitmap) {
235        let (columns, vis) = self.data.into_parts();
236        (self.ops, columns, vis)
237    }
238
239    pub fn to_protobuf(&self) -> PbStreamChunk {
240        if !self.is_vis_compacted() {
241            return self.clone().compact_vis().to_protobuf();
242        }
243        PbStreamChunk {
244            cardinality: self.cardinality() as u32,
245            ops: self.ops.iter().map(|op| op.to_protobuf() as i32).collect(),
246            columns: self.columns().iter().map(|col| col.to_protobuf()).collect(),
247        }
248    }
249
250    pub fn from_protobuf(prost: &PbStreamChunk) -> ArrayResult<Self> {
251        let cardinality = prost.get_cardinality() as usize;
252        let mut ops = Vec::with_capacity(cardinality);
253        for op in prost.get_ops() {
254            ops.push(Op::from_protobuf(op)?);
255        }
256        let mut columns = vec![];
257        for column in prost.get_columns() {
258            columns.push(ArrayImpl::from_protobuf(column, cardinality)?.into());
259        }
260        Ok(StreamChunk::new(ops, columns))
261    }
262
263    pub fn ops(&self) -> &[Op] {
264        &self.ops
265    }
266
267    /// Returns a table-like text representation of the `StreamChunk`.
268    pub fn to_pretty(&self) -> impl Display + use<> {
269        self.to_pretty_inner(None)
270    }
271
272    /// Returns a table-like text representation of the `StreamChunk` with a header of column names
273    /// from the given `schema`.
274    pub fn to_pretty_with_schema(&self, schema: &Schema) -> impl Display + use<> {
275        self.to_pretty_inner(Some(schema))
276    }
277
278    fn to_pretty_inner(&self, schema: Option<&Schema>) -> impl Display + use<> {
279        use comfy_table::{Cell, CellAlignment, Table};
280
281        if self.cardinality() == 0 {
282            return Either::Left("(empty)");
283        }
284
285        let mut table = Table::new();
286        table.load_preset(DataChunk::PRETTY_TABLE_PRESET);
287
288        if let Some(schema) = schema {
289            assert_eq!(self.dimension(), schema.len());
290            let cells = std::iter::once(String::new())
291                .chain(schema.fields().iter().map(|f| f.name.clone()));
292            table.set_header(cells);
293        }
294
295        for (op, row_ref) in self.rows() {
296            let mut cells = Vec::with_capacity(row_ref.len() + 1);
297            cells.push(
298                Cell::new(match op {
299                    Op::Insert => "+",
300                    Op::Delete => "-",
301                    Op::UpdateDelete => "U-",
302                    Op::UpdateInsert => "U+",
303                })
304                .set_alignment(CellAlignment::Right),
305            );
306            for datum in row_ref.iter() {
307                let str = match datum {
308                    None => "".to_owned(), // NULL
309                    Some(scalar) => scalar.to_text(),
310                };
311                cells.push(Cell::new(str));
312            }
313            table.add_row(cells);
314        }
315
316        Either::Right(table)
317    }
318
319    /// Reorder (and possibly remove) columns.
320    ///
321    /// e.g. if `indices` is `[2, 1, 0]`, and the chunk contains column `[a, b, c]`, then the output
322    /// will be `[c, b, a]`. If `indices` is [2, 0], then the output will be `[c, a]`.
323    /// If the input mapping is identity mapping, no reorder will be performed.
324    pub fn project(&self, indices: &[usize]) -> Self {
325        Self {
326            ops: self.ops.clone(),
327            data: self.data.project(indices),
328        }
329    }
330
331    /// Remove the adjacent delete-insert and insert-deletes if their row value are the same.
332    pub fn eliminate_adjacent_noop_update(self) -> Self {
333        let len = self.data_chunk().capacity();
334        let mut c: StreamChunkMut = self.into();
335        let mut prev_r = None;
336        for curr in 0..len {
337            if !c.vis(curr) {
338                continue;
339            }
340            if let Some(prev) = prev_r
341                && (
342                    // 1. Delete then Insert
343                    (matches!(c.op(prev), Op::UpdateDelete | Op::Delete)
344                    && matches!(c.op(curr), Op::UpdateInsert | Op::Insert))
345                    ||
346                    // 2. Insert then Delete
347                    //
348                    // Note that after eliminating `U+` and `U-` here, we will get a new
349                    // pair of `U-` and `U+` consisting of `prev.prev` and the `next` row.
350                    // `prev.prev` and `prev`, `curr` and `next` share the same stream key
351                    // because they are `U-` and `U+` pairs, while `prev` and `curr` share
352                    // the same stream key because they are equal, so `prev-prev` and `next`
353                    // must also share the same stream key. Therefore, it's okay to leave
354                    // their `Update` ops unchanged.
355                    //
356                    // Also, they can't be the same, otherwise these 4 rows are the same,
357                    // which should already be eliminated by the first branch.
358                    (matches!(c.op(prev), Op::UpdateInsert | Op::Insert)
359                    && matches!(c.op(curr), Op::UpdateDelete | Op::Delete))
360                )
361                && c.row_ref(prev) == c.row_ref(curr)
362            {
363                c.set_vis(prev, false);
364                c.set_vis(curr, false);
365                prev_r = None;
366            } else {
367                prev_r = Some(curr);
368            }
369        }
370
371        // Normalize update pairs that became partially invisible.
372        // If only U- is visible, turn it into Delete; if only U+ is visible, turn it into Insert.
373        for idx in 0..len.saturating_sub(1) {
374            if c.op(idx) == Op::UpdateDelete && c.op(idx + 1) == Op::UpdateInsert {
375                let delete_vis = c.vis(idx);
376                let insert_vis = c.vis(idx + 1);
377                if delete_vis && !insert_vis {
378                    c.set_op(idx, Op::Delete);
379                } else if !delete_vis && insert_vis {
380                    c.set_op(idx + 1, Op::Insert);
381                }
382            }
383        }
384        c.into()
385    }
386
387    /// Reorder columns and set visibility.
388    pub fn project_with_vis(&self, indices: &[usize], vis: Bitmap) -> Self {
389        Self {
390            ops: self.ops.clone(),
391            data: self.data.project_with_vis(indices, vis),
392        }
393    }
394
395    /// Clone the `StreamChunk` with a new visibility.
396    pub fn clone_with_vis(&self, vis: Bitmap) -> Self {
397        Self {
398            ops: self.ops.clone(),
399            data: self.data.with_visibility(vis),
400        }
401    }
402
403    /// Keeps visible rows whose operation is in `ops`
404    pub fn retain_ops(self, ops: &[Op]) -> Self {
405        let mut chunk: StreamChunkMut = self.into();
406
407        for (_, mut row) in chunk.to_rows_mut() {
408            if !ops.contains(&row.op()) {
409                row.set_vis(false);
410            }
411        }
412
413        chunk.into()
414    }
415}
416
417impl Deref for StreamChunk {
418    type Target = DataChunk;
419
420    fn deref(&self) -> &Self::Target {
421        &self.data
422    }
423}
424
425impl DerefMut for StreamChunk {
426    fn deref_mut(&mut self) -> &mut Self::Target {
427        &mut self.data
428    }
429}
430
431/// `StreamChunk` can be created from `DataChunk` with all operations set to `Insert`.
432impl From<DataChunk> for StreamChunk {
433    fn from(data: DataChunk) -> Self {
434        Self::from_parts(vec![Op::Insert; data.capacity()], data)
435    }
436}
437
438impl fmt::Debug for StreamChunk {
439    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
440        if f.alternate() {
441            write!(
442                f,
443                "StreamChunk {{ cardinality: {}, capacity: {}, data:\n{}\n }}",
444                self.cardinality(),
445                self.capacity(),
446                self.to_pretty()
447            )
448        } else {
449            f.debug_struct("StreamChunk")
450                .field("cardinality", &self.cardinality())
451                .field("capacity", &self.capacity())
452                .finish_non_exhaustive()
453        }
454    }
455}
456
457impl EstimateSize for StreamChunk {
458    fn estimated_heap_size(&self) -> usize {
459        self.data.estimated_heap_size() + self.ops.len() * size_of::<Op>()
460    }
461}
462
463enum OpsMutState {
464    ArcRef(Arc<[Op]>),
465    Mut(Vec<Op>),
466}
467
468impl OpsMutState {
469    const UNDEFINED: Self = Self::Mut(Vec::new());
470}
471
472pub struct OpsMut {
473    state: OpsMutState,
474}
475
476impl OpsMut {
477    pub fn new(ops: Arc<[Op]>) -> Self {
478        Self {
479            state: OpsMutState::ArcRef(ops),
480        }
481    }
482
483    pub fn len(&self) -> usize {
484        match &self.state {
485            OpsMutState::ArcRef(v) => v.len(),
486            OpsMutState::Mut(v) => v.len(),
487        }
488    }
489
490    pub fn is_empty(&self) -> bool {
491        self.len() == 0
492    }
493
494    pub fn set(&mut self, n: usize, val: Op) {
495        debug_assert!(n < self.len());
496        if let OpsMutState::Mut(v) = &mut self.state {
497            v[n] = val;
498        } else {
499            let state = mem::replace(&mut self.state, OpsMutState::UNDEFINED); // intermediate state
500            let mut v = match state {
501                OpsMutState::ArcRef(v) => v.to_vec(),
502                OpsMutState::Mut(_) => unreachable!(),
503            };
504            v[n] = val;
505            self.state = OpsMutState::Mut(v);
506        }
507    }
508
509    pub fn get(&self, n: usize) -> Op {
510        debug_assert!(n < self.len());
511        match &self.state {
512            OpsMutState::ArcRef(v) => v[n],
513            OpsMutState::Mut(v) => v[n],
514        }
515    }
516}
517impl From<OpsMut> for Arc<[Op]> {
518    fn from(v: OpsMut) -> Self {
519        match v.state {
520            OpsMutState::ArcRef(a) => a,
521            OpsMutState::Mut(v) => v.into(),
522        }
523    }
524}
525
526/// A mutable wrapper for `StreamChunk`. can only set the visibilities and ops in place, can not
527/// change the length.
528pub struct StreamChunkMut {
529    columns: Arc<[ArrayRef]>,
530    ops: OpsMut,
531    vis: BitmapBuilder,
532}
533
534impl From<StreamChunk> for StreamChunkMut {
535    fn from(c: StreamChunk) -> Self {
536        let (c, ops) = c.into_parts();
537        let (columns, vis) = c.into_parts_v2();
538        Self {
539            columns,
540            ops: OpsMut::new(ops),
541            vis: vis.into(),
542        }
543    }
544}
545
546impl From<StreamChunkMut> for StreamChunk {
547    fn from(c: StreamChunkMut) -> Self {
548        StreamChunk::from_parts(c.ops, DataChunk::from_parts(c.columns, c.vis.finish()))
549    }
550}
551
552/// A handle to one row of a [`StreamChunkMut`] that can update the row's op and
553/// visibility in place. It holds a raw pointer instead of `&mut` since multiple
554/// handles to the same chunk coexist (see [`StreamChunkMut::to_rows_mut`]).
555pub struct OpRowMutRef<'a> {
556    c: *mut StreamChunkMut,
557    i: usize,
558    _phantom: PhantomData<&'a mut StreamChunkMut>,
559}
560
561impl PartialEq for OpRowMutRef<'_> {
562    fn eq(&self, other: &Self) -> bool {
563        self.row_ref() == other.row_ref()
564    }
565}
566impl Eq for OpRowMutRef<'_> {}
567
568impl<'a> OpRowMutRef<'a> {
569    pub fn index(&self) -> usize {
570        self.i
571    }
572
573    // SAFETY of derefs below: `self.c` is valid for `'a`, and each access reborrows
574    // a single disjoint field only.
575
576    pub fn vis(&self) -> bool {
577        let vis = unsafe { &(*self.c).vis };
578        vis.is_set(self.i)
579    }
580
581    pub fn op(&self) -> Op {
582        let ops = unsafe { &(*self.c).ops };
583        ops.get(self.i)
584    }
585
586    pub fn set_vis(&mut self, val: bool) {
587        let vis = unsafe { &mut (*self.c).vis };
588        vis.set(self.i, val);
589    }
590
591    pub fn set_op(&mut self, val: Op) {
592        let ops = unsafe { &mut (*self.c).ops };
593        ops.set(self.i, val);
594    }
595
596    pub fn row_ref(&self) -> RowRef<'_> {
597        RowRef::with_columns(unsafe { &(*self.c).columns }, self.i)
598    }
599
600    /// return if the two row ref is in the same chunk
601    pub fn same_chunk(&self, other: &Self) -> bool {
602        std::ptr::eq(self.c, other.c)
603    }
604}
605
606impl StreamChunkMut {
607    pub fn capacity(&self) -> usize {
608        self.vis.len()
609    }
610
611    pub fn vis(&self, i: usize) -> bool {
612        self.vis.is_set(i)
613    }
614
615    pub fn op(&self, i: usize) -> Op {
616        self.ops.get(i)
617    }
618
619    pub fn row_ref(&self, i: usize) -> RowRef<'_> {
620        RowRef::with_columns(self.columns(), i)
621    }
622
623    pub fn set_vis(&mut self, n: usize, val: bool) {
624        self.vis.set(n, val);
625    }
626
627    pub fn set_op(&mut self, n: usize, val: Op) {
628        self.ops.set(n, val);
629    }
630
631    pub fn columns(&self) -> &[ArrayRef] {
632        &self.columns
633    }
634
635    /// get the mut reference of the stream chunk.
636    pub fn to_rows_mut(&mut self) -> impl Iterator<Item = (RowRef<'_>, OpRowMutRef<'_>)> {
637        // SAFETY: the pointer is derived from `&mut self`, which stays exclusively
638        // borrowed by the returned iterator.
639        unsafe { Self::rows_mut_ptr(self) }
640    }
641
642    /// # Safety
643    ///
644    /// `p` must be derived from an exclusive reference valid for `'a`, and the chunk must
645    /// not be accessed in other ways while the iterator or any yielded item is alive.
646    unsafe fn rows_mut_ptr<'a>(
647        p: *mut Self,
648    ) -> impl Iterator<Item = (RowRef<'a>, OpRowMutRef<'a>)> {
649        let len = {
650            let vis = unsafe { &(*p).vis };
651            vis.len()
652        };
653        (0..len)
654            .filter(move |i| {
655                let vis = unsafe { &(*p).vis };
656                vis.is_set(*i)
657            })
658            .map(move |i| {
659                (
660                    RowRef::with_columns(unsafe { &(*p).columns }, i),
661                    OpRowMutRef {
662                        c: p,
663                        i,
664                        _phantom: PhantomData,
665                    },
666                )
667            })
668    }
669}
670
671/// Test utilities for [`StreamChunk`].
672#[easy_ext::ext(StreamChunkTestExt)]
673impl StreamChunk {
674    /// Parse a chunk from string.
675    ///
676    /// See also [`DataChunkTestExt::from_pretty`].
677    ///
678    /// # Format
679    ///
680    /// The first line is a header indicating the column types.
681    /// The following lines indicate rows within the chunk.
682    /// Each line starts with an operation followed by values.
683    /// NULL values are represented as `.`.
684    ///
685    /// # Example
686    /// ```
687    /// use risingwave_common::array::StreamChunk;
688    /// use risingwave_common::array::stream_chunk::StreamChunkTestExt as _;
689    /// let chunk = StreamChunk::from_pretty(
690    ///     "  I I I I      // type chars
691    ///     U- 2 5 . .      // '.' means NULL
692    ///     U+ 2 5 2 6 D    // 'D' means deleted in visibility
693    ///     +  . . 4 8      // ^ comments are ignored
694    ///     -  . . 3 4",
695    /// );
696    /// //  ^ operations:
697    /// //     +: Insert
698    /// //     -: Delete
699    /// //    U+: UpdateInsert
700    /// //    U-: UpdateDelete
701    ///
702    /// // type chars:
703    /// //     I: i64
704    /// //     i: i32
705    /// //     F: f64
706    /// //     f: f32
707    /// //     T: str
708    /// //    TS: Timestamp
709    /// //    TZ: Timestamptz
710    /// //   SRL: Serial
711    /// //   x[]: array of x
712    /// // <i,f>: struct
713    /// ```
714    pub fn from_pretty(s: &str) -> Self {
715        let mut chunk_str = String::new();
716        let mut ops = vec![];
717
718        let (header, body) = match s.split_once('\n') {
719            Some(pair) => pair,
720            None => {
721                // empty chunk
722                return StreamChunk {
723                    ops: Arc::new([]),
724                    data: DataChunk::from_pretty(s),
725                };
726            }
727        };
728        chunk_str.push_str(header);
729        chunk_str.push('\n');
730
731        for line in body.split_inclusive('\n') {
732            if line.trim_start().is_empty() {
733                continue;
734            }
735            let (op, row) = line
736                .trim_start()
737                .split_once(|c: char| c.is_ascii_whitespace())
738                .ok_or_else(|| panic!("missing operation: {line:?}"))
739                .unwrap();
740            ops.push(match op {
741                "+" => Op::Insert,
742                "-" => Op::Delete,
743                "U+" => Op::UpdateInsert,
744                "U-" => Op::UpdateDelete,
745                t => panic!("invalid op: {t:?}"),
746            });
747            chunk_str.push_str(row);
748        }
749        StreamChunk {
750            ops: ops.into(),
751            data: DataChunk::from_pretty(&chunk_str),
752        }
753    }
754
755    /// Validate the `StreamChunk` layout.
756    pub fn valid(&self) -> bool {
757        let len = self.ops.len();
758        let data = &self.data;
759        data.visibility().len() == len && data.columns().iter().all(|col| col.len() == len)
760    }
761
762    /// Concatenate multiple `StreamChunk` into one.
763    ///
764    /// Panics if `chunks` is empty.
765    pub fn concat(chunks: Vec<StreamChunk>) -> StreamChunk {
766        let data_types = chunks[0].data_types();
767        let size = chunks.iter().map(|c| c.cardinality()).sum::<usize>();
768
769        let mut builder = StreamChunkBuilder::unlimited(data_types, Some(size));
770
771        for chunk in chunks {
772            // TODO: directly append chunks.
773            for (op, row) in chunk.rows() {
774                let none = builder.append_row(op, row);
775                debug_assert!(none.is_none());
776            }
777        }
778
779        builder.take().expect("chunk should not be empty")
780    }
781
782    /// Sort rows.
783    pub fn sort_rows(self) -> Self {
784        if self.capacity() == 0 {
785            return self;
786        }
787        let rows = self.rows().collect_vec();
788        let mut idx = (0..self.capacity()).collect_vec();
789        idx.sort_by_key(|&i| {
790            let (op, row_ref) = rows[i];
791            (op, DefaultOrdered(row_ref))
792        });
793        StreamChunk {
794            ops: idx.iter().map(|&i| self.ops[i]).collect(),
795            data: self.data.reorder_rows(&idx),
796        }
797    }
798
799    /// Generate `num_of_chunks` data chunks with type `data_types`,
800    /// where each data chunk has cardinality of `chunk_size`.
801    /// TODO(kwannoel): Generate different types of op, different vis.
802    pub fn gen_stream_chunks(
803        num_of_chunks: usize,
804        chunk_size: usize,
805        data_types: &[DataType],
806        varchar_properties: &VarcharProperty,
807    ) -> Vec<StreamChunk> {
808        Self::gen_stream_chunks_inner(
809            num_of_chunks,
810            chunk_size,
811            data_types,
812            varchar_properties,
813            1.0,
814            1.0,
815        )
816    }
817
818    pub fn gen_stream_chunks_inner(
819        num_of_chunks: usize,
820        chunk_size: usize,
821        data_types: &[DataType],
822        varchar_properties: &VarcharProperty,
823        visibility_percent: f64, // % of rows that are visible
824        inserts_percent: f64,    // Rest will be deletes.
825    ) -> Vec<StreamChunk> {
826        let ops = if inserts_percent == 0.0 {
827            vec![Op::Delete; chunk_size]
828        } else if inserts_percent == 1.0 {
829            vec![Op::Insert; chunk_size]
830        } else {
831            let mut rng = SmallRng::from_seed([0; 32]);
832            let mut ops = vec![];
833            for _ in 0..chunk_size {
834                ops.push(if rng.random_bool(inserts_percent) {
835                    Op::Insert
836                } else {
837                    Op::Delete
838                });
839            }
840            ops
841        };
842        DataChunk::gen_data_chunks(
843            num_of_chunks,
844            chunk_size,
845            data_types,
846            varchar_properties,
847            visibility_percent,
848        )
849        .into_iter()
850        .map(|chunk| StreamChunk::from_parts(ops.clone(), chunk))
851        .collect()
852    }
853}
854
855#[cfg(test)]
856mod tests {
857    use super::*;
858
859    #[test]
860    fn test_to_pretty_string() {
861        let chunk = StreamChunk::from_pretty(
862            "  I I
863             + 1 6
864             - 2 .
865            U- 3 7
866            U+ 4 .",
867        );
868        assert_eq!(
869            chunk.to_pretty().to_string(),
870            "\
871+----+---+---+
872|  + | 1 | 6 |
873|  - | 2 |   |
874| U- | 3 | 7 |
875| U+ | 4 |   |
876+----+---+---+"
877        );
878    }
879
880    #[test]
881    fn test_split_1() {
882        let chunk = StreamChunk::from_pretty(
883            "  I I
884             + 1 6
885             - 2 .
886            U- 3 7
887            U+ 4 .",
888        );
889        let results = chunk.split(2);
890        assert_eq!(2, results.len());
891        assert_eq!(
892            results[0].to_pretty().to_string(),
893            "\
894+---+---+---+
895| + | 1 | 6 |
896| - | 2 |   |
897+---+---+---+"
898        );
899        assert_eq!(
900            results[1].to_pretty().to_string(),
901            "\
902+----+---+---+
903| U- | 3 | 7 |
904| U+ | 4 |   |
905+----+---+---+"
906        );
907    }
908
909    #[test]
910    fn test_split_2() {
911        let chunk = StreamChunk::from_pretty(
912            "  I I
913             + 1 6
914            U- 3 7
915            U+ 4 .
916             - 2 .",
917        );
918        let results = chunk.split(2);
919        assert_eq!(2, results.len());
920        assert_eq!(
921            results[0].to_pretty().to_string(),
922            "\
923+----+---+---+
924|  + | 1 | 6 |
925| U- | 3 | 7 |
926| U+ | 4 |   |
927+----+---+---+"
928        );
929        assert_eq!(
930            results[1].to_pretty().to_string(),
931            "\
932+---+---+---+
933| - | 2 |   |
934+---+---+---+"
935        );
936    }
937
938    #[test]
939    fn test_eliminate_adjacent_noop_update() {
940        let c = StreamChunk::from_pretty(
941            "  I I
942            - 1 6 D
943            - 2 2
944            + 2 3
945            - 2 3
946            + 1 6
947            - 1 7
948            + 1 10 D
949            + 1 7
950            U- 3 7
951            U+ 3 7
952            + 2 3",
953        );
954        let c = c.eliminate_adjacent_noop_update();
955        assert_eq!(
956            c.to_pretty().to_string(),
957            "\
958+---+---+---+
959| - | 2 | 2 |
960| + | 1 | 6 |
961| + | 2 | 3 |
962+---+---+---+"
963        );
964    }
965
966    #[test]
967    fn test_eliminate_adjacent_noop_update_normalize_update_pair() {
968        let c = StreamChunk::from_pretty(
969            "  I I
970            + 1 10
971            U- 1 10
972            U+ 1 20",
973        );
974        let c = c.eliminate_adjacent_noop_update();
975        assert_eq!(
976            c.to_pretty().to_string(),
977            "\
978+---+---+----+
979| + | 1 | 20 |
980+---+---+----+"
981        );
982    }
983}