Skip to main content

risingwave_stream/executor/iceberg_with_pk_index/
position_delete_merger.rs

1// Copyright 2026 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 anyhow::Context;
16use risingwave_common::id::SinkId;
17use risingwave_connector::sink::Result as SinkResult;
18use risingwave_pb::connector_service::SinkMetadata;
19use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
20
21use crate::executor::prelude::*;
22use crate::task::LocalBarrierManager;
23
24/// Trait abstracting position-delete file operations for testability.
25///
26/// Implementations are responsible for reading existing position deletes
27/// (V3 Puffin deletion vectors or V2 Parquet position-delete files),
28/// merging new delete positions, writing the resulting delete file, and
29/// returning the commit metadata for the current barrier.
30#[async_trait::async_trait]
31pub trait PositionDeleteHandler: Send + 'static {
32    fn start_seed(&mut self, wait_epoch: u64);
33
34    fn write(&mut self, path: &str, pos: i64) -> SinkResult<()>;
35
36    async fn flush(&mut self) -> SinkResult<Option<SinkMetadata>>;
37}
38
39/// Position-delete merger executor for iceberg pk-index sink without Equality Delete.
40///
41/// This stateless executor receives [`file_path`, `position`] messages from the Writer Executor,
42/// merges them with existing position deletes, and reports the merged delete-file metadata to
43/// meta on each barrier. Depending on the table format version, the written delete file is a
44/// V3 Puffin deletion vector or a V2 file-scoped Parquet position-delete file.
45///
46/// The upstream plan shards messages by `file_path`, so each actor only merges delete
47/// positions for the files assigned to its shard.
48///
49/// Input schema: [`file_path`: Varchar, `position`: int64]
50/// Output: Barriers and watermarks only; no data chunks (terminal executor in the stream graph).
51pub struct PositionDeleteMergerExecutor<H>
52where
53    H: PositionDeleteHandler,
54{
55    actor_id: ActorId,
56    sink_id: SinkId,
57    local_barrier_manager: LocalBarrierManager,
58    input: Option<Executor>,
59    handler: H,
60}
61
62impl<H> PositionDeleteMergerExecutor<H>
63where
64    H: PositionDeleteHandler,
65{
66    pub fn new(
67        actor_id: ActorId,
68        sink_id: SinkId,
69        local_barrier_manager: LocalBarrierManager,
70        input: Executor,
71        handler: H,
72    ) -> Self {
73        Self {
74            actor_id,
75            sink_id,
76            local_barrier_manager,
77            input: Some(input),
78            handler,
79        }
80    }
81
82    #[try_stream(ok = Message, error = StreamExecutorError)]
83    async fn execute_inner(mut self) {
84        let mut input = self.input.take().unwrap().execute();
85
86        // Consume the first barrier. Its `prev` epoch is the point the previous actors (if this is a
87        // scale) committed up to; seed only from a snapshot that includes it.
88        let barrier = expect_first_barrier(&mut input).await?;
89        self.handler.start_seed(barrier.epoch.prev);
90        yield Message::Barrier(barrier);
91
92        #[for_await]
93        for msg in input {
94            match msg? {
95                Message::Chunk(chunk) => {
96                    for (op, row) in chunk.rows() {
97                        debug_assert_eq!(op, risingwave_common::array::Op::Insert);
98                        let file_path = row
99                            .datum_at(0)
100                            .map(|d| d.into_utf8())
101                            .context("file_path should not be null")?;
102                        let position = row
103                            .datum_at(1)
104                            .context("position should not be null")?
105                            .into_int64();
106                        self.handler
107                            .write(file_path, position)
108                            .map_err(|e| StreamExecutorError::sink_error(e, self.sink_id))?;
109                    }
110                }
111                Message::Barrier(barrier) => {
112                    barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
113
114                    let mut metadata = None;
115                    if barrier.is_checkpoint() {
116                        metadata = self
117                            .handler
118                            .flush()
119                            .await
120                            .map_err(|e| StreamExecutorError::sink_error(e, self.sink_id))?;
121                    }
122
123                    if let Some(metadata) = metadata
124                        && metadata.metadata.is_some()
125                    {
126                        self.local_barrier_manager
127                            .report_iceberg_pk_index_sink_metadata(
128                                barrier.epoch,
129                                self.sink_id,
130                                self.actor_id,
131                                PbIcebergPkIndexSinkRole::PositionDeleteMerger,
132                                Some(metadata),
133                            );
134                    }
135
136                    yield Message::Barrier(barrier);
137                }
138                Message::Watermark(w) => {
139                    yield Message::Watermark(w);
140                }
141            }
142        }
143    }
144}
145
146impl<H> Execute for PositionDeleteMergerExecutor<H>
147where
148    H: PositionDeleteHandler,
149{
150    fn execute(self: Box<Self>) -> BoxedMessageStream {
151        self.execute_inner().boxed()
152    }
153}
154
155#[cfg(test)]
156mod tests {
157    use std::collections::BTreeSet;
158    use std::sync::{Arc, Mutex};
159
160    use hashbrown::HashMap;
161    use risingwave_common::array::{Array, ArrayBuilder, I64ArrayBuilder, Op, Utf8ArrayBuilder};
162    use risingwave_common::catalog::{Field, Schema};
163    use risingwave_common::id::SinkId;
164    use risingwave_common::types::DataType;
165    use risingwave_common::util::epoch::test_epoch;
166
167    use super::*;
168    use crate::executor::test_utils::MockSource;
169    use crate::task::LocalBarrierManager;
170
171    fn build_delete_position_chunk(positions: &[(&str, i64)]) -> StreamChunk {
172        let len = positions.len();
173        let mut file_path_builder = Utf8ArrayBuilder::new(len);
174        let mut position_builder = I64ArrayBuilder::new(len);
175
176        for (path, offset) in positions {
177            file_path_builder.append(Some(*path));
178            position_builder.append(Some(*offset));
179        }
180
181        StreamChunk::from_parts(
182            vec![Op::Insert; len],
183            risingwave_common::array::DataChunk::new(
184                vec![
185                    file_path_builder.finish().into_ref(),
186                    position_builder.finish().into_ref(),
187                ],
188                len,
189            ),
190        )
191    }
192
193    #[derive(Clone)]
194    struct PositionDeleteHandlerMock {
195        existing_dvs: Arc<Mutex<HashMap<String, BTreeSet<i64>>>>,
196        pending_dvs: Arc<Mutex<HashMap<String, BTreeSet<i64>>>>,
197        written_dvs: Arc<Mutex<HashMap<String, BTreeSet<i64>>>>,
198    }
199
200    impl PositionDeleteHandlerMock {
201        fn new() -> Self {
202            Self {
203                existing_dvs: Arc::new(Mutex::new(HashMap::new())),
204                pending_dvs: Arc::new(Mutex::new(HashMap::new())),
205                written_dvs: Arc::new(Mutex::new(HashMap::new())),
206            }
207        }
208
209        fn with_existing_dv(self, file_path: &str, positions: BTreeSet<i64>) -> Self {
210            self.existing_dvs
211                .lock()
212                .unwrap()
213                .insert(file_path.to_owned(), positions);
214            self
215        }
216    }
217
218    #[async_trait::async_trait]
219    impl PositionDeleteHandler for PositionDeleteHandlerMock {
220        fn start_seed(&mut self, _wait_epoch: u64) {}
221
222        fn write(&mut self, path: &str, pos: i64) -> SinkResult<()> {
223            self.pending_dvs
224                .lock()
225                .unwrap()
226                .entry_ref(path)
227                .or_default()
228                .insert(pos);
229            Ok(())
230        }
231
232        async fn flush(&mut self) -> SinkResult<Option<SinkMetadata>> {
233            let pending = {
234                let mut pending_dvs = self.pending_dvs.lock().unwrap();
235                std::mem::take(&mut *pending_dvs)
236            };
237            if pending.is_empty() {
238                return Ok(None);
239            }
240
241            let mut existing_dvs = self.existing_dvs.lock().unwrap();
242            let mut written_dvs = self.written_dvs.lock().unwrap();
243            for (file_path, positions) in pending {
244                let mut merged = existing_dvs.get(&file_path).cloned().unwrap_or_default();
245                merged.extend(positions);
246                existing_dvs.insert(file_path.clone(), merged.clone());
247                written_dvs.insert(file_path, merged);
248            }
249            // The mock asserts on side effects (`written_dvs`) rather than emitted
250            // metadata, so we still return Ok(None). The report-on-barrier path
251            // is exercised by the SLT integration tests.
252            // TODO: add unit test for report-on-barrier path once test infra is
253            // available to capture `LocalBarrierEvent`s on the receiver side.
254            Ok(None)
255        }
256    }
257
258    fn input_schema() -> Schema {
259        Schema::new(vec![
260            Field::unnamed(DataType::Varchar),
261            Field::unnamed(DataType::Int64),
262        ])
263    }
264
265    #[tokio::test]
266    async fn test_position_delete_merger_basic() {
267        let handler = PositionDeleteHandlerMock::new();
268        let written_dvs = handler.written_dvs.clone();
269
270        let (mut tx, source) = MockSource::channel();
271        let source = source.into_executor(input_schema(), vec![]);
272
273        let lbm = LocalBarrierManager::for_test();
274        let mut executor =
275            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
276                .boxed()
277                .execute();
278
279        tx.push_barrier(test_epoch(1), false);
280        assert!(executor.next().await.unwrap().unwrap().is_barrier());
281
282        tx.push_chunk(build_delete_position_chunk(&[
283            ("file1.parquet", 0),
284            ("file1.parquet", 3),
285            ("file2.parquet", 1),
286        ]));
287        tx.push_barrier(test_epoch(2), false);
288
289        assert!(executor.next().await.unwrap().unwrap().is_barrier());
290
291        let dvs = written_dvs.lock().unwrap();
292        assert_eq!(dvs.get("file1.parquet").unwrap(), &BTreeSet::from([0, 3]));
293        assert_eq!(dvs.get("file2.parquet").unwrap(), &BTreeSet::from([1]));
294    }
295
296    #[tokio::test]
297    async fn test_position_delete_merger_merge_with_existing() {
298        let handler = PositionDeleteHandlerMock::new()
299            .with_existing_dv("file1.parquet", BTreeSet::from([0, 5, 10]));
300        let written_dvs = handler.written_dvs.clone();
301
302        let (mut tx, source) = MockSource::channel();
303        let source = source.into_executor(input_schema(), vec![]);
304
305        let lbm = LocalBarrierManager::for_test();
306        let mut executor =
307            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
308                .boxed()
309                .execute();
310
311        tx.push_barrier(test_epoch(1), false);
312        assert!(executor.next().await.unwrap().unwrap().is_barrier());
313
314        tx.push_chunk(build_delete_position_chunk(&[
315            ("file1.parquet", 3),
316            ("file1.parquet", 5),
317            ("file1.parquet", 7),
318        ]));
319        tx.push_barrier(test_epoch(2), false);
320
321        assert!(executor.next().await.unwrap().unwrap().is_barrier());
322
323        let dvs = written_dvs.lock().unwrap();
324        assert_eq!(
325            dvs.get("file1.parquet").unwrap(),
326            &BTreeSet::from([0, 3, 5, 7, 10])
327        );
328    }
329
330    #[tokio::test]
331    async fn test_position_delete_merger_no_deletes() {
332        let handler = PositionDeleteHandlerMock::new();
333        let written_dvs = handler.written_dvs.clone();
334
335        let (mut tx, source) = MockSource::channel();
336        let source = source.into_executor(input_schema(), vec![]);
337
338        let lbm = LocalBarrierManager::for_test();
339        let mut executor =
340            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
341                .boxed()
342                .execute();
343
344        tx.push_barrier(test_epoch(1), false);
345        assert!(executor.next().await.unwrap().unwrap().is_barrier());
346
347        tx.push_barrier(test_epoch(2), false);
348        assert!(executor.next().await.unwrap().unwrap().is_barrier());
349
350        assert!(written_dvs.lock().unwrap().is_empty());
351    }
352
353    #[tokio::test]
354    async fn test_position_delete_merger_multiple_epochs() {
355        let handler = PositionDeleteHandlerMock::new();
356        let written_dvs = handler.written_dvs.clone();
357
358        let (mut tx, source) = MockSource::channel();
359        let source = source.into_executor(input_schema(), vec![]);
360
361        let lbm = LocalBarrierManager::for_test();
362        let mut executor =
363            PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
364                .boxed()
365                .execute();
366
367        tx.push_barrier(test_epoch(1), false);
368        assert!(executor.next().await.unwrap().unwrap().is_barrier());
369
370        tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 0)]));
371        tx.push_barrier(test_epoch(2), false);
372        assert!(executor.next().await.unwrap().unwrap().is_barrier());
373        assert_eq!(
374            written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
375            &BTreeSet::from([0])
376        );
377
378        tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 2)]));
379        tx.push_barrier(test_epoch(3), false);
380        assert!(executor.next().await.unwrap().unwrap().is_barrier());
381
382        assert_eq!(
383            written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
384            &BTreeSet::from([0, 2])
385        );
386    }
387}