1use anyhow::Context;
16use risingwave_common::bail;
17use risingwave_common::id::SinkId;
18use risingwave_connector::sink::Result as SinkResult;
19use risingwave_pb::connector_service::SinkMetadata;
20use risingwave_pb::id::IcebergCompactionTaskId;
21use risingwave_pb::stream_plan::iceberg_pk_index_compaction_context::Phase;
22use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
23
24use crate::executor::prelude::*;
25use crate::task::LocalBarrierManager;
26
27#[async_trait::async_trait]
34pub trait PositionDeleteHandler: Send + 'static {
35 fn start_seed(&mut self, wait_epoch: u64);
42
43 fn write(&mut self, path: &str, pos: i64) -> SinkResult<()>;
44
45 async fn flush(&mut self) -> SinkResult<Option<SinkMetadata>>;
46}
47
48pub struct PositionDeleteMergerExecutor<H>
67where
68 H: PositionDeleteHandler,
69{
70 actor_id: ActorId,
71 sink_id: SinkId,
72 local_barrier_manager: LocalBarrierManager,
73 input: Option<Executor>,
74 handler: H,
75}
76
77impl<H> PositionDeleteMergerExecutor<H>
78where
79 H: PositionDeleteHandler,
80{
81 pub fn new(
82 actor_id: ActorId,
83 sink_id: SinkId,
84 local_barrier_manager: LocalBarrierManager,
85 input: Executor,
86 handler: H,
87 ) -> Self {
88 Self {
89 actor_id,
90 sink_id,
91 local_barrier_manager,
92 input: Some(input),
93 handler,
94 }
95 }
96
97 #[try_stream(ok = Message, error = StreamExecutorError)]
98 async fn execute_inner(mut self) {
99 let mut input = self.input.take().unwrap().execute();
100
101 let barrier = expect_first_barrier(&mut input).await?;
104 self.handler.start_seed(barrier.epoch.prev);
105 yield Message::Barrier(barrier);
106
107 #[for_await]
108 for msg in input {
109 match msg? {
110 Message::Chunk(chunk) => {
111 for (op, row) in chunk.rows() {
112 debug_assert_eq!(op, risingwave_common::array::Op::Insert);
113 let file_path = row
114 .datum_at(0)
115 .map(|d| d.into_utf8())
116 .context("file_path should not be null")?;
117 let position = row
118 .datum_at(1)
119 .context("position should not be null")?
120 .into_int64();
121 self.handler
122 .write(file_path, position)
123 .map_err(|e| StreamExecutorError::sink_error(e, self.sink_id))?;
124 }
125 }
126 Message::Barrier(barrier) => {
127 barrier.assume_no_update_vnode_bitmap(self.actor_id)?;
128 let compaction_resumed = self.compaction_resume_task_id(&barrier);
129 if compaction_resumed.is_some() && !barrier.is_checkpoint() {
130 bail!(
131 "iceberg pk-index merger {} received a non-checkpoint compaction resume barrier {:?}",
132 self.sink_id,
133 barrier.epoch
134 );
135 }
136
137 let mut metadata = None;
138 if barrier.is_checkpoint() {
139 metadata = self
140 .handler
141 .flush()
142 .await
143 .map_err(|e| StreamExecutorError::sink_error(e, self.sink_id))?;
144 }
145
146 if let Some(metadata) = metadata
147 && metadata.metadata.is_some()
148 {
149 self.local_barrier_manager
150 .report_iceberg_pk_index_sink_metadata(
151 barrier.epoch,
152 self.sink_id,
153 self.actor_id,
154 PbIcebergPkIndexSinkRole::PositionDeleteMerger,
155 Some(metadata),
156 );
157 }
158
159 let reseed = compaction_resumed.map(|_| barrier.epoch.prev);
160 yield Message::Barrier(barrier);
161
162 if let Some(epoch) = reseed {
163 self.handler.start_seed(epoch);
164 }
165 }
166 Message::Watermark(w) => {
167 yield Message::Watermark(w);
168 }
169 }
170 }
171 }
172
173 fn compaction_resume_task_id(&self, barrier: &Barrier) -> Option<IcebergCompactionTaskId> {
176 match barrier.iceberg_pk_index_compaction() {
177 Some(context)
178 if context.sink_id == self.sink_id && context.phase == Phase::End as i32 =>
179 {
180 Some(context.task_id)
181 }
182 _ => None,
183 }
184 }
185}
186
187impl<H> Execute for PositionDeleteMergerExecutor<H>
188where
189 H: PositionDeleteHandler,
190{
191 fn execute(self: Box<Self>) -> BoxedMessageStream {
192 self.execute_inner().boxed()
193 }
194}
195
196#[cfg(test)]
197mod tests {
198 use std::collections::BTreeSet;
199 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
200 use std::sync::{Arc, Mutex};
201
202 use hashbrown::HashMap;
203 use risingwave_common::array::{Array, ArrayBuilder, I64ArrayBuilder, Op, Utf8ArrayBuilder};
204 use risingwave_common::catalog::{Field, Schema};
205 use risingwave_common::id::SinkId;
206 use risingwave_common::types::DataType;
207 use risingwave_common::util::epoch::{EpochPair, test_epoch};
208
209 use super::*;
210 use crate::executor::test_utils::MockSource;
211 use crate::task::LocalBarrierManager;
212
213 fn build_delete_position_chunk(positions: &[(&str, i64)]) -> StreamChunk {
214 let len = positions.len();
215 let mut file_path_builder = Utf8ArrayBuilder::new(len);
216 let mut position_builder = I64ArrayBuilder::new(len);
217
218 for (path, offset) in positions {
219 file_path_builder.append(Some(*path));
220 position_builder.append(Some(*offset));
221 }
222
223 StreamChunk::from_parts(
224 vec![Op::Insert; len],
225 risingwave_common::array::DataChunk::new(
226 vec![
227 file_path_builder.finish().into_ref(),
228 position_builder.finish().into_ref(),
229 ],
230 len,
231 ),
232 )
233 }
234
235 type Dvs = HashMap<String, BTreeSet<i64>>;
236
237 #[derive(Clone)]
247 struct PositionDeleteHandlerMock {
248 table_state: Arc<Mutex<Dvs>>,
249 staged: Arc<Mutex<Dvs>>,
250 pending: Arc<Mutex<Dvs>>,
251 written_dvs: Arc<Mutex<Dvs>>,
252 seed_epochs: Arc<Mutex<Vec<u64>>>,
254 block_second_seed: Arc<AtomicBool>,
255 second_seed_started: Arc<AtomicBool>,
256 release_second_seed: Arc<tokio::sync::Notify>,
257 flush_calls: Arc<AtomicUsize>,
258 }
259
260 impl PositionDeleteHandlerMock {
261 fn new() -> Self {
262 Self {
263 table_state: Arc::new(Mutex::new(HashMap::new())),
264 staged: Arc::new(Mutex::new(HashMap::new())),
265 pending: Arc::new(Mutex::new(HashMap::new())),
266 written_dvs: Arc::new(Mutex::new(HashMap::new())),
267 seed_epochs: Arc::new(Mutex::new(Vec::new())),
268 block_second_seed: Arc::new(AtomicBool::new(false)),
269 second_seed_started: Arc::new(AtomicBool::new(false)),
270 release_second_seed: Arc::new(tokio::sync::Notify::new()),
271 flush_calls: Arc::new(AtomicUsize::new(0)),
272 }
273 }
274
275 fn with_existing_dv(self, file_path: &str, positions: BTreeSet<i64>) -> Self {
278 self.table_state
279 .lock()
280 .unwrap()
281 .insert(file_path.to_owned(), positions);
282 self
283 }
284 }
285
286 #[async_trait::async_trait]
287 impl PositionDeleteHandler for PositionDeleteHandlerMock {
288 fn start_seed(&mut self, wait_epoch: u64) {
289 let mut seed_epochs = self.seed_epochs.lock().unwrap();
290 seed_epochs.push(wait_epoch);
291 let is_second_seed = seed_epochs.len() == 2;
292 drop(seed_epochs);
293 if is_second_seed && self.block_second_seed.load(Ordering::SeqCst) {
294 self.second_seed_started.store(true, Ordering::SeqCst);
295 return;
296 }
297 *self.staged.lock().unwrap() = self.table_state.lock().unwrap().clone();
301 }
302
303 fn write(&mut self, path: &str, pos: i64) -> SinkResult<()> {
304 self.pending
305 .lock()
306 .unwrap()
307 .entry_ref(path)
308 .or_default()
309 .insert(pos);
310 Ok(())
311 }
312
313 async fn flush(&mut self) -> SinkResult<Option<SinkMetadata>> {
314 self.flush_calls.fetch_add(1, Ordering::SeqCst);
315 if self.second_seed_started.load(Ordering::SeqCst)
316 && self.block_second_seed.load(Ordering::SeqCst)
317 {
318 self.release_second_seed.notified().await;
319 *self.staged.lock().unwrap() = self.table_state.lock().unwrap().clone();
320 }
321 let pending = std::mem::take(&mut *self.pending.lock().unwrap());
322 if pending.is_empty() {
323 return Ok(None);
324 }
325
326 let mut staged = self.staged.lock().unwrap();
327 let mut table_state = self.table_state.lock().unwrap();
328 let mut written_dvs = self.written_dvs.lock().unwrap();
329 for (file_path, positions) in pending {
330 let mut merged = staged.get(&file_path).cloned().unwrap_or_default();
333 merged.extend(positions);
334 staged.insert(file_path.clone(), merged.clone());
335 table_state.insert(file_path.clone(), merged.clone());
336 written_dvs.insert(file_path, merged);
337 }
338 Ok(None)
344 }
345 }
346
347 fn compaction_barrier(epoch: u64, sink_id: SinkId, phase: Phase) -> Barrier {
348 Barrier::new_test_barrier(test_epoch(epoch)).with_iceberg_pk_index_compaction(
349 crate::executor::IcebergPkIndexCompactionContext {
350 sink_id,
351 task_id: 7.into(),
352 phase: phase as i32,
353 },
354 )
355 }
356
357 const MERGER_ACTOR_ID: ActorId = ActorId::new(123);
358
359 fn input_schema() -> Schema {
360 Schema::new(vec![
361 Field::unnamed(DataType::Varchar),
362 Field::unnamed(DataType::Int64),
363 ])
364 }
365
366 #[tokio::test]
367 async fn test_position_delete_merger_basic() {
368 let handler = PositionDeleteHandlerMock::new();
369 let written_dvs = handler.written_dvs.clone();
370
371 let (mut tx, source) = MockSource::channel();
372 let source = source.into_executor(input_schema(), vec![]);
373
374 let lbm = LocalBarrierManager::for_test();
375 let mut executor =
376 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
377 .boxed()
378 .execute();
379
380 tx.push_barrier(test_epoch(1), false);
381 assert!(executor.next().await.unwrap().unwrap().is_barrier());
382
383 tx.push_chunk(build_delete_position_chunk(&[
384 ("file1.parquet", 0),
385 ("file1.parquet", 3),
386 ("file2.parquet", 1),
387 ]));
388 tx.push_barrier(test_epoch(2), false);
389
390 assert!(executor.next().await.unwrap().unwrap().is_barrier());
391
392 let dvs = written_dvs.lock().unwrap();
393 assert_eq!(dvs.get("file1.parquet").unwrap(), &BTreeSet::from([0, 3]));
394 assert_eq!(dvs.get("file2.parquet").unwrap(), &BTreeSet::from([1]));
395 }
396
397 #[tokio::test]
398 async fn test_position_delete_merger_merge_with_existing() {
399 let handler = PositionDeleteHandlerMock::new()
400 .with_existing_dv("file1.parquet", BTreeSet::from([0, 5, 10]));
401 let written_dvs = handler.written_dvs.clone();
402
403 let (mut tx, source) = MockSource::channel();
404 let source = source.into_executor(input_schema(), vec![]);
405
406 let lbm = LocalBarrierManager::for_test();
407 let mut executor =
408 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
409 .boxed()
410 .execute();
411
412 tx.push_barrier(test_epoch(1), false);
413 assert!(executor.next().await.unwrap().unwrap().is_barrier());
414
415 tx.push_chunk(build_delete_position_chunk(&[
416 ("file1.parquet", 3),
417 ("file1.parquet", 5),
418 ("file1.parquet", 7),
419 ]));
420 tx.push_barrier(test_epoch(2), false);
421
422 assert!(executor.next().await.unwrap().unwrap().is_barrier());
423
424 let dvs = written_dvs.lock().unwrap();
425 assert_eq!(
426 dvs.get("file1.parquet").unwrap(),
427 &BTreeSet::from([0, 3, 5, 7, 10])
428 );
429 }
430
431 #[tokio::test]
432 async fn test_position_delete_merger_no_deletes() {
433 let handler = PositionDeleteHandlerMock::new();
434 let written_dvs = handler.written_dvs.clone();
435
436 let (mut tx, source) = MockSource::channel();
437 let source = source.into_executor(input_schema(), vec![]);
438
439 let lbm = LocalBarrierManager::for_test();
440 let mut executor =
441 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
442 .boxed()
443 .execute();
444
445 tx.push_barrier(test_epoch(1), false);
446 assert!(executor.next().await.unwrap().unwrap().is_barrier());
447
448 tx.push_barrier(test_epoch(2), false);
449 assert!(executor.next().await.unwrap().unwrap().is_barrier());
450
451 assert!(written_dvs.lock().unwrap().is_empty());
452 }
453
454 #[tokio::test]
455 async fn test_position_delete_merger_multiple_epochs() {
456 let handler = PositionDeleteHandlerMock::new();
457 let written_dvs = handler.written_dvs.clone();
458
459 let (mut tx, source) = MockSource::channel();
460 let source = source.into_executor(input_schema(), vec![]);
461
462 let lbm = LocalBarrierManager::for_test();
463 let mut executor =
464 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
465 .boxed()
466 .execute();
467
468 tx.push_barrier(test_epoch(1), false);
469 assert!(executor.next().await.unwrap().unwrap().is_barrier());
470
471 tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 0)]));
472 tx.push_barrier(test_epoch(2), false);
473 assert!(executor.next().await.unwrap().unwrap().is_barrier());
474 assert_eq!(
475 written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
476 &BTreeSet::from([0])
477 );
478
479 tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 2)]));
480 tx.push_barrier(test_epoch(3), false);
481 assert!(executor.next().await.unwrap().unwrap().is_barrier());
482
483 assert_eq!(
484 written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
485 &BTreeSet::from([0, 2])
486 );
487 }
488
489 #[tokio::test]
490 async fn test_compaction_b2_forwards_before_reseed_wait() {
491 let sink_id = SinkId::new(0);
492 let handler = PositionDeleteHandlerMock::new();
493 let seed_epochs = handler.seed_epochs.clone();
494
495 let (mut tx, source) = MockSource::channel();
496 let source = source.into_executor(input_schema(), vec![]);
497 let mut executor = PositionDeleteMergerExecutor::new(
498 MERGER_ACTOR_ID,
499 sink_id,
500 LocalBarrierManager::for_test(),
501 source,
502 handler,
503 )
504 .boxed()
505 .execute();
506
507 tx.push_barrier(test_epoch(1), false);
508 assert!(executor.next().await.unwrap().unwrap().is_barrier());
509 let initial_prev = EpochPair::new_test_epoch(test_epoch(1)).prev;
510
511 tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
512 tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
513 let b2 = executor.next().await.unwrap().unwrap();
514 assert!(b2.is_barrier(), "B2 must be observable before re-seeding");
515 assert_eq!(
516 *seed_epochs.lock().unwrap(),
517 vec![initial_prev],
518 "the generator must not call start_seed until resumed after yielding B2"
519 );
520
521 tx.push_barrier(test_epoch(3), false);
522 assert!(executor.next().await.unwrap().unwrap().is_barrier());
523 assert_eq!(
524 *seed_epochs.lock().unwrap(),
525 vec![initial_prev, EpochPair::new_test_epoch(test_epoch(2)).prev]
526 );
527 }
528
529 #[tokio::test]
530 async fn test_compaction_reseed_buffers_e2_until_next_checkpoint() {
531 let sink_id = SinkId::new(0);
532 let handler = PositionDeleteHandlerMock::new();
533 handler.block_second_seed.store(true, Ordering::SeqCst);
534 let second_seed_started = handler.second_seed_started.clone();
535 let release_second_seed = handler.release_second_seed.clone();
536 let flush_calls = handler.flush_calls.clone();
537 let written_dvs = handler.written_dvs.clone();
538
539 let (mut tx, source) = MockSource::channel();
540 let source = source.into_executor(input_schema(), vec![]);
541 let mut executor = PositionDeleteMergerExecutor::new(
542 MERGER_ACTOR_ID,
543 sink_id,
544 LocalBarrierManager::for_test(),
545 source,
546 handler,
547 )
548 .boxed()
549 .execute();
550
551 tx.push_barrier(test_epoch(1), false);
552 assert!(executor.next().await.unwrap().unwrap().is_barrier());
553 tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
554 assert!(executor.next().await.unwrap().unwrap().is_barrier());
555 assert!(
556 !second_seed_started.load(Ordering::SeqCst),
557 "B2 must be forwarded before the re-seed starts"
558 );
559
560 tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
561 tx.push_barrier(test_epoch(3), false);
562 assert!(
563 tokio::time::timeout(std::time::Duration::from_millis(20), executor.next())
564 .await
565 .is_err(),
566 "the next checkpoint must wait for the post-B2 seed"
567 );
568 assert!(second_seed_started.load(Ordering::SeqCst));
569 assert_eq!(flush_calls.load(Ordering::SeqCst), 2);
570 assert!(
571 written_dvs.lock().unwrap().is_empty(),
572 "E2 deletes must remain buffered until seeding completes"
573 );
574
575 release_second_seed.notify_one();
576 assert!(executor.next().await.unwrap().unwrap().is_barrier());
577 assert_eq!(
578 written_dvs.lock().unwrap().get("output.parquet").unwrap(),
579 &BTreeSet::from([9])
580 );
581 }
582
583 #[tokio::test]
588 async fn test_position_delete_merger_compaction_resume_reseeds_from_new_snapshot() {
589 let sink_id = SinkId::new(0);
590 let handler = PositionDeleteHandlerMock::new();
591 let table_state = handler.table_state.clone();
592 let staged = handler.staged.clone();
593 let written_dvs = handler.written_dvs.clone();
594 let seed_epochs = handler.seed_epochs.clone();
595
596 let (mut tx, source) = MockSource::channel();
597 let source = source.into_executor(input_schema(), vec![]);
598
599 let lbm = LocalBarrierManager::for_test();
600 let mut executor =
601 PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
602 .boxed()
603 .execute();
604
605 tx.push_barrier(test_epoch(1), false);
606 assert!(executor.next().await.unwrap().unwrap().is_barrier());
607
608 tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
610 tx.push_barrier(test_epoch(2), false);
611 assert!(executor.next().await.unwrap().unwrap().is_barrier());
612 assert_eq!(
613 written_dvs.lock().unwrap().get("input.parquet").unwrap(),
614 &BTreeSet::from([0])
615 );
616
617 {
620 let mut table_state = table_state.lock().unwrap();
621 table_state.remove("input.parquet");
622 table_state.insert("output.parquet".to_owned(), BTreeSet::from([4]));
623 }
624
625 tx.send_barrier(compaction_barrier(3, sink_id, Phase::End));
626 assert!(executor.next().await.unwrap().unwrap().is_barrier());
627 assert_eq!(
628 *seed_epochs.lock().unwrap(),
629 vec![EpochPair::new_test_epoch(test_epoch(1)).prev],
630 "re-seed starts only after B2 is forwarded"
631 );
632
633 tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
634 tx.push_barrier(test_epoch(4), false);
635 assert!(executor.next().await.unwrap().unwrap().is_barrier());
636
637 {
638 let staged = staged.lock().unwrap();
639 assert!(
640 !staged.contains_key("input.parquet"),
641 "the retired data file must be evicted from the resident cache"
642 );
643 assert_eq!(
644 staged.get("output.parquet").unwrap(),
645 &BTreeSet::from([4, 9]),
646 "the resolver's delete and E2 delete must be resident after the re-seed and flush"
647 );
648 }
649
650 assert_eq!(
651 written_dvs.lock().unwrap().get("output.parquet").unwrap(),
652 &BTreeSet::from([4, 9])
653 );
654
655 assert_eq!(
657 *seed_epochs.lock().unwrap(),
658 vec![
659 EpochPair::new_test_epoch(test_epoch(1)).prev,
660 EpochPair::new_test_epoch(test_epoch(3)).prev,
661 ]
662 );
663 }
664
665 #[tokio::test]
668 async fn test_position_delete_merger_compaction_begin_does_not_reseed() {
669 let sink_id = SinkId::new(0);
670 let handler = PositionDeleteHandlerMock::new();
671 let seed_epochs = handler.seed_epochs.clone();
672
673 let (mut tx, source) = MockSource::channel();
674 let source = source.into_executor(input_schema(), vec![]);
675
676 let lbm = LocalBarrierManager::for_test();
677 let mut executor =
678 PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
679 .boxed()
680 .execute();
681
682 tx.push_barrier(test_epoch(1), false);
683 assert!(executor.next().await.unwrap().unwrap().is_barrier());
684
685 tx.send_barrier(compaction_barrier(2, sink_id, Phase::Begin));
686 assert!(executor.next().await.unwrap().unwrap().is_barrier());
687
688 assert_eq!(seed_epochs.lock().unwrap().len(), 1);
689 }
690
691 #[tokio::test]
694 async fn test_position_delete_merger_compaction_resume_for_other_sink_is_ignored() {
695 let handler = PositionDeleteHandlerMock::new();
696 let seed_epochs = handler.seed_epochs.clone();
697
698 let (mut tx, source) = MockSource::channel();
699 let source = source.into_executor(input_schema(), vec![]);
700
701 let lbm = LocalBarrierManager::for_test();
702 let mut executor = PositionDeleteMergerExecutor::new(
703 MERGER_ACTOR_ID,
704 SinkId::new(0),
705 lbm,
706 source,
707 handler,
708 )
709 .boxed()
710 .execute();
711
712 tx.push_barrier(test_epoch(1), false);
713 assert!(executor.next().await.unwrap().unwrap().is_barrier());
714
715 tx.send_barrier(compaction_barrier(2, SinkId::new(1), Phase::End));
716 assert!(executor.next().await.unwrap().unwrap().is_barrier());
717
718 assert_eq!(seed_epochs.lock().unwrap().len(), 1);
719 }
720}