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 risingwave_pb::stream_plan::IcebergPkIndexCompactionContext {
350 sink_id,
351 task_id: 7.into(),
352 phase: phase as i32,
353 resolver_task_input: (phase == Phase::Begin).then_some(
354 risingwave_pb::stream_plan::iceberg_pk_index_compaction_context::ResolverTaskInput::default(),
355 ),
356 },
357 )
358 }
359
360 const MERGER_ACTOR_ID: ActorId = ActorId::new(123);
361
362 fn input_schema() -> Schema {
363 Schema::new(vec![
364 Field::unnamed(DataType::Varchar),
365 Field::unnamed(DataType::Int64),
366 ])
367 }
368
369 #[tokio::test]
370 async fn test_position_delete_merger_basic() {
371 let handler = PositionDeleteHandlerMock::new();
372 let written_dvs = handler.written_dvs.clone();
373
374 let (mut tx, source) = MockSource::channel();
375 let source = source.into_executor(input_schema(), vec![]);
376
377 let lbm = LocalBarrierManager::for_test();
378 let mut executor =
379 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
380 .boxed()
381 .execute();
382
383 tx.push_barrier(test_epoch(1), false);
384 assert!(executor.next().await.unwrap().unwrap().is_barrier());
385
386 tx.push_chunk(build_delete_position_chunk(&[
387 ("file1.parquet", 0),
388 ("file1.parquet", 3),
389 ("file2.parquet", 1),
390 ]));
391 tx.push_barrier(test_epoch(2), false);
392
393 assert!(executor.next().await.unwrap().unwrap().is_barrier());
394
395 let dvs = written_dvs.lock().unwrap();
396 assert_eq!(dvs.get("file1.parquet").unwrap(), &BTreeSet::from([0, 3]));
397 assert_eq!(dvs.get("file2.parquet").unwrap(), &BTreeSet::from([1]));
398 }
399
400 #[tokio::test]
401 async fn test_position_delete_merger_merge_with_existing() {
402 let handler = PositionDeleteHandlerMock::new()
403 .with_existing_dv("file1.parquet", BTreeSet::from([0, 5, 10]));
404 let written_dvs = handler.written_dvs.clone();
405
406 let (mut tx, source) = MockSource::channel();
407 let source = source.into_executor(input_schema(), vec![]);
408
409 let lbm = LocalBarrierManager::for_test();
410 let mut executor =
411 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
412 .boxed()
413 .execute();
414
415 tx.push_barrier(test_epoch(1), false);
416 assert!(executor.next().await.unwrap().unwrap().is_barrier());
417
418 tx.push_chunk(build_delete_position_chunk(&[
419 ("file1.parquet", 3),
420 ("file1.parquet", 5),
421 ("file1.parquet", 7),
422 ]));
423 tx.push_barrier(test_epoch(2), false);
424
425 assert!(executor.next().await.unwrap().unwrap().is_barrier());
426
427 let dvs = written_dvs.lock().unwrap();
428 assert_eq!(
429 dvs.get("file1.parquet").unwrap(),
430 &BTreeSet::from([0, 3, 5, 7, 10])
431 );
432 }
433
434 #[tokio::test]
435 async fn test_position_delete_merger_no_deletes() {
436 let handler = PositionDeleteHandlerMock::new();
437 let written_dvs = handler.written_dvs.clone();
438
439 let (mut tx, source) = MockSource::channel();
440 let source = source.into_executor(input_schema(), vec![]);
441
442 let lbm = LocalBarrierManager::for_test();
443 let mut executor =
444 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
445 .boxed()
446 .execute();
447
448 tx.push_barrier(test_epoch(1), false);
449 assert!(executor.next().await.unwrap().unwrap().is_barrier());
450
451 tx.push_barrier(test_epoch(2), false);
452 assert!(executor.next().await.unwrap().unwrap().is_barrier());
453
454 assert!(written_dvs.lock().unwrap().is_empty());
455 }
456
457 #[tokio::test]
458 async fn test_position_delete_merger_multiple_epochs() {
459 let handler = PositionDeleteHandlerMock::new();
460 let written_dvs = handler.written_dvs.clone();
461
462 let (mut tx, source) = MockSource::channel();
463 let source = source.into_executor(input_schema(), vec![]);
464
465 let lbm = LocalBarrierManager::for_test();
466 let mut executor =
467 PositionDeleteMergerExecutor::new(123.into(), SinkId::new(0), lbm, source, handler)
468 .boxed()
469 .execute();
470
471 tx.push_barrier(test_epoch(1), false);
472 assert!(executor.next().await.unwrap().unwrap().is_barrier());
473
474 tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 0)]));
475 tx.push_barrier(test_epoch(2), false);
476 assert!(executor.next().await.unwrap().unwrap().is_barrier());
477 assert_eq!(
478 written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
479 &BTreeSet::from([0])
480 );
481
482 tx.push_chunk(build_delete_position_chunk(&[("file1.parquet", 2)]));
483 tx.push_barrier(test_epoch(3), false);
484 assert!(executor.next().await.unwrap().unwrap().is_barrier());
485
486 assert_eq!(
487 written_dvs.lock().unwrap().get("file1.parquet").unwrap(),
488 &BTreeSet::from([0, 2])
489 );
490 }
491
492 #[tokio::test]
493 async fn test_compaction_b2_forwards_before_reseed_wait() {
494 let sink_id = SinkId::new(0);
495 let handler = PositionDeleteHandlerMock::new();
496 let seed_epochs = handler.seed_epochs.clone();
497
498 let (mut tx, source) = MockSource::channel();
499 let source = source.into_executor(input_schema(), vec![]);
500 let mut executor = PositionDeleteMergerExecutor::new(
501 MERGER_ACTOR_ID,
502 sink_id,
503 LocalBarrierManager::for_test(),
504 source,
505 handler,
506 )
507 .boxed()
508 .execute();
509
510 tx.push_barrier(test_epoch(1), false);
511 assert!(executor.next().await.unwrap().unwrap().is_barrier());
512 let initial_prev = EpochPair::new_test_epoch(test_epoch(1)).prev;
513
514 tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
515 tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
516 let b2 = executor.next().await.unwrap().unwrap();
517 assert!(b2.is_barrier(), "B2 must be observable before re-seeding");
518 assert_eq!(
519 *seed_epochs.lock().unwrap(),
520 vec![initial_prev],
521 "the generator must not call start_seed until resumed after yielding B2"
522 );
523
524 tx.push_barrier(test_epoch(3), false);
525 assert!(executor.next().await.unwrap().unwrap().is_barrier());
526 assert_eq!(
527 *seed_epochs.lock().unwrap(),
528 vec![initial_prev, EpochPair::new_test_epoch(test_epoch(2)).prev]
529 );
530 }
531
532 #[tokio::test]
533 async fn test_compaction_reseed_buffers_e2_until_next_checkpoint() {
534 let sink_id = SinkId::new(0);
535 let handler = PositionDeleteHandlerMock::new();
536 handler.block_second_seed.store(true, Ordering::SeqCst);
537 let second_seed_started = handler.second_seed_started.clone();
538 let release_second_seed = handler.release_second_seed.clone();
539 let flush_calls = handler.flush_calls.clone();
540 let written_dvs = handler.written_dvs.clone();
541
542 let (mut tx, source) = MockSource::channel();
543 let source = source.into_executor(input_schema(), vec![]);
544 let mut executor = PositionDeleteMergerExecutor::new(
545 MERGER_ACTOR_ID,
546 sink_id,
547 LocalBarrierManager::for_test(),
548 source,
549 handler,
550 )
551 .boxed()
552 .execute();
553
554 tx.push_barrier(test_epoch(1), false);
555 assert!(executor.next().await.unwrap().unwrap().is_barrier());
556 tx.send_barrier(compaction_barrier(2, sink_id, Phase::End));
557 assert!(executor.next().await.unwrap().unwrap().is_barrier());
558 assert!(
559 !second_seed_started.load(Ordering::SeqCst),
560 "B2 must be forwarded before the re-seed starts"
561 );
562
563 tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
564 tx.push_barrier(test_epoch(3), false);
565 assert!(
566 tokio::time::timeout(std::time::Duration::from_millis(20), executor.next())
567 .await
568 .is_err(),
569 "the next checkpoint must wait for the post-B2 seed"
570 );
571 assert!(second_seed_started.load(Ordering::SeqCst));
572 assert_eq!(flush_calls.load(Ordering::SeqCst), 2);
573 assert!(
574 written_dvs.lock().unwrap().is_empty(),
575 "E2 deletes must remain buffered until seeding completes"
576 );
577
578 release_second_seed.notify_one();
579 assert!(executor.next().await.unwrap().unwrap().is_barrier());
580 assert_eq!(
581 written_dvs.lock().unwrap().get("output.parquet").unwrap(),
582 &BTreeSet::from([9])
583 );
584 }
585
586 #[tokio::test]
591 async fn test_position_delete_merger_compaction_resume_reseeds_from_new_snapshot() {
592 let sink_id = SinkId::new(0);
593 let handler = PositionDeleteHandlerMock::new();
594 let table_state = handler.table_state.clone();
595 let staged = handler.staged.clone();
596 let written_dvs = handler.written_dvs.clone();
597 let seed_epochs = handler.seed_epochs.clone();
598
599 let (mut tx, source) = MockSource::channel();
600 let source = source.into_executor(input_schema(), vec![]);
601
602 let lbm = LocalBarrierManager::for_test();
603 let mut executor =
604 PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
605 .boxed()
606 .execute();
607
608 tx.push_barrier(test_epoch(1), false);
609 assert!(executor.next().await.unwrap().unwrap().is_barrier());
610
611 tx.push_chunk(build_delete_position_chunk(&[("input.parquet", 0)]));
613 tx.push_barrier(test_epoch(2), false);
614 assert!(executor.next().await.unwrap().unwrap().is_barrier());
615 assert_eq!(
616 written_dvs.lock().unwrap().get("input.parquet").unwrap(),
617 &BTreeSet::from([0])
618 );
619
620 {
623 let mut table_state = table_state.lock().unwrap();
624 table_state.remove("input.parquet");
625 table_state.insert("output.parquet".to_owned(), BTreeSet::from([4]));
626 }
627
628 tx.send_barrier(compaction_barrier(3, sink_id, Phase::End));
629 assert!(executor.next().await.unwrap().unwrap().is_barrier());
630 assert_eq!(
631 *seed_epochs.lock().unwrap(),
632 vec![EpochPair::new_test_epoch(test_epoch(1)).prev],
633 "re-seed starts only after B2 is forwarded"
634 );
635
636 tx.push_chunk(build_delete_position_chunk(&[("output.parquet", 9)]));
637 tx.push_barrier(test_epoch(4), false);
638 assert!(executor.next().await.unwrap().unwrap().is_barrier());
639
640 {
641 let staged = staged.lock().unwrap();
642 assert!(
643 !staged.contains_key("input.parquet"),
644 "the retired data file must be evicted from the resident cache"
645 );
646 assert_eq!(
647 staged.get("output.parquet").unwrap(),
648 &BTreeSet::from([4, 9]),
649 "the resolver's delete and E2 delete must be resident after the re-seed and flush"
650 );
651 }
652
653 assert_eq!(
654 written_dvs.lock().unwrap().get("output.parquet").unwrap(),
655 &BTreeSet::from([4, 9])
656 );
657
658 assert_eq!(
660 *seed_epochs.lock().unwrap(),
661 vec![
662 EpochPair::new_test_epoch(test_epoch(1)).prev,
663 EpochPair::new_test_epoch(test_epoch(3)).prev,
664 ]
665 );
666 }
667
668 #[tokio::test]
671 async fn test_position_delete_merger_compaction_begin_does_not_reseed() {
672 let sink_id = SinkId::new(0);
673 let handler = PositionDeleteHandlerMock::new();
674 let seed_epochs = handler.seed_epochs.clone();
675
676 let (mut tx, source) = MockSource::channel();
677 let source = source.into_executor(input_schema(), vec![]);
678
679 let lbm = LocalBarrierManager::for_test();
680 let mut executor =
681 PositionDeleteMergerExecutor::new(MERGER_ACTOR_ID, sink_id, lbm, source, handler)
682 .boxed()
683 .execute();
684
685 tx.push_barrier(test_epoch(1), false);
686 assert!(executor.next().await.unwrap().unwrap().is_barrier());
687
688 tx.send_barrier(compaction_barrier(2, sink_id, Phase::Begin));
689 assert!(executor.next().await.unwrap().unwrap().is_barrier());
690
691 assert_eq!(seed_epochs.lock().unwrap().len(), 1);
692 }
693
694 #[tokio::test]
697 async fn test_position_delete_merger_compaction_resume_for_other_sink_is_ignored() {
698 let handler = PositionDeleteHandlerMock::new();
699 let seed_epochs = handler.seed_epochs.clone();
700
701 let (mut tx, source) = MockSource::channel();
702 let source = source.into_executor(input_schema(), vec![]);
703
704 let lbm = LocalBarrierManager::for_test();
705 let mut executor = PositionDeleteMergerExecutor::new(
706 MERGER_ACTOR_ID,
707 SinkId::new(0),
708 lbm,
709 source,
710 handler,
711 )
712 .boxed()
713 .execute();
714
715 tx.push_barrier(test_epoch(1), false);
716 assert!(executor.next().await.unwrap().unwrap().is_barrier());
717
718 tx.send_barrier(compaction_barrier(2, SinkId::new(1), Phase::End));
719 assert!(executor.next().await.unwrap().unwrap().is_barrier());
720
721 assert_eq!(seed_epochs.lock().unwrap().len(), 1);
722 }
723}