risingwave_stream/executor/iceberg_with_pk_index/
position_delete_merger.rs1use 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#[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
39pub 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 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 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}