1use anyhow::Context;
16use iceberg::writer::PositionDeleteInput;
17use risingwave_common::array::stream_record::Record;
18use risingwave_common::array::{DataChunk, Op};
19use risingwave_common::bail;
20use risingwave_common::id::SinkId;
21use risingwave_common::row::{Project, RowExt};
22use risingwave_common::util::chunk_coalesce::DataChunkBuilder;
23use risingwave_common::util::epoch::EpochPair;
24use risingwave_common::util::iter_util::ZipEqFast;
25use risingwave_pb::connector_service::SinkMetadata;
26use risingwave_pb::id::IcebergCompactionTaskId;
27use risingwave_pb::stream_plan::iceberg_pk_index_compaction_context::Phase;
28use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
29use risingwave_storage::StateStore;
30
31use crate::common::change_buffer::output_kind;
32use crate::common::compact_chunk::{InconsistencyBehavior, compact_chunk_inline};
33use crate::executor::prelude::*;
34use crate::task::LocalBarrierManager;
35
36type PkRow<'a> = Project<'a, RowRef<'a>>;
37
38fn new_chunk_builder(chunk_size: usize) -> DataChunkBuilder {
39 DataChunkBuilder::new(vec![DataType::Varchar, DataType::Int64], chunk_size)
40}
41
42fn append_row(builder: &mut DataChunkBuilder, file_path: &str, position: i64) -> Option<DataChunk> {
43 builder.append_one_row([
44 Some(ScalarRefImpl::Utf8(file_path)),
45 Some(ScalarRefImpl::Int64(position)),
46 ])
47}
48
49#[derive(Debug)]
51pub enum WriterInputMode {
52 Normal,
53 ResolvingRight {
54 task_id: IcebergCompactionTaskId,
55 begin_epoch: EpochPair,
56 },
57 DrainingLeft {
58 task_id: IcebergCompactionTaskId,
59 barrier: Barrier,
60 },
61}
62
63#[async_trait::async_trait]
68pub trait IcebergWriter: Send + 'static {
69 async fn write_chunk(
71 &mut self,
72 chunk: DataChunk,
73 ) -> StreamExecutorResult<Vec<PositionDeleteInput>>;
74
75 async fn flush(&mut self) -> StreamExecutorResult<Option<SinkMetadata>>;
78}
79
80pub struct WriterExecutor<S, W>
94where
95 S: StateStore,
96 W: IcebergWriter,
97{
98 ctx: ActorContextRef,
99 input: Option<Executor>,
100 resolver_input: Option<Executor>,
101 pk_indices: Vec<usize>,
103 pk_index_state_table: StateTable<S>,
106 writer: W,
108 delete_position_buffer: Option<DataChunkBuilder>,
110 mode: WriterInputMode,
111 chunk_size: usize,
112 sink_id: SinkId,
113 local_barrier_manager: LocalBarrierManager,
114}
115
116impl<S, W> WriterExecutor<S, W>
117where
118 S: StateStore,
119 W: IcebergWriter,
120{
121 #[expect(clippy::too_many_arguments)]
122 pub fn new(
123 ctx: ActorContextRef,
124 input: Executor,
125 resolver_input: Executor,
126 pk_indices: Vec<usize>,
127 pk_index_state_table: StateTable<S>,
128 writer: W,
129 chunk_size: usize,
130 sink_id: SinkId,
131 local_barrier_manager: LocalBarrierManager,
132 ) -> Self {
133 Self {
134 ctx,
135 input: Some(input),
136 resolver_input: Some(resolver_input),
137 pk_indices,
138 pk_index_state_table,
139 writer,
140 delete_position_buffer: None,
141 mode: WriterInputMode::Normal,
142 chunk_size,
143 sink_id,
144 local_barrier_manager,
145 }
146 }
147
148 async fn delete_existing_row(
149 &mut self,
150 pk_row: PkRow<'_>,
151 delete_position_buffer: &mut DataChunkBuilder,
152 ) -> StreamExecutorResult<Option<DataChunk>> {
153 let Some(index_row) = self.pk_index_state_table.get_row(pk_row).await? else {
154 return Ok(None);
155 };
156
157 let num_cols = index_row.len();
158 let file_path = index_row
159 .datum_at(num_cols - 2)
160 .context("file_path should not be null")?
161 .into_utf8();
162 let position = index_row
163 .datum_at(num_cols - 1)
164 .context("position should not be null")?
165 .into_int64();
166 let chunk = append_row(delete_position_buffer, file_path, position);
167 self.pk_index_state_table.delete(index_row);
168 Ok(chunk)
169 }
170
171 #[try_stream(ok = DataChunk, error = StreamExecutorError)]
186 async fn process_chunk(&mut self, chunk: StreamChunk) {
187 let chunk = compact_chunk_inline::<{ output_kind::RETRACT }>(
188 chunk,
189 &self.pk_indices,
190 InconsistencyBehavior::Panic,
191 );
192
193 let mut delete_position_buffer = self
194 .delete_position_buffer
195 .take()
196 .unwrap_or_else(|| new_chunk_builder(self.chunk_size));
197 let pk_indices = self.pk_indices.clone();
198
199 let mut insert_chunk =
207 DataChunkBuilder::new(chunk.data_chunk().data_types(), chunk.capacity() + 1);
208 let mut insert_pks: Vec<PkRow<'_>> = Vec::new();
209
210 for record in chunk.records() {
211 match record {
212 Record::Insert { new_row } => {
213 let overflow = insert_chunk.append_one_row(new_row);
214 debug_assert!(overflow.is_none(), "insert chunk exceeds capacity");
215 insert_pks.push(new_row.project(&pk_indices));
216 }
217 Record::Delete { old_row } => {
218 let pk_row = old_row.project(&pk_indices);
219 if let Some(chunk) = self
220 .delete_existing_row(pk_row, &mut delete_position_buffer)
221 .await?
222 {
223 yield chunk;
224 }
225 }
226 Record::Update { new_row, .. } => {
227 let pk_row = new_row.project(&pk_indices);
229 if let Some(chunk) = self
230 .delete_existing_row(pk_row, &mut delete_position_buffer)
231 .await?
232 {
233 yield chunk;
234 }
235 let overflow = insert_chunk.append_one_row(new_row);
236 debug_assert!(overflow.is_none(), "insert chunk exceeds capacity");
237 insert_pks.push(pk_row);
238 }
239 }
240 }
241
242 if !insert_pks.is_empty() {
243 let write_chunk = insert_chunk.finish();
244 let positions = self.writer.write_chunk(write_chunk).await?;
245
246 for (pk, pos) in insert_pks.into_iter().zip_eq_fast(positions) {
247 let mut index_row_data = Vec::with_capacity(pk_indices.len() + 2);
248 for datum in pk.iter() {
249 index_row_data.push(datum);
250 }
251 index_row_data.push(Some(ScalarRefImpl::Utf8(&pos.path)));
252 index_row_data.push(Some(ScalarRefImpl::Int64(pos.pos)));
253 self.pk_index_state_table.insert(index_row_data.as_slice());
254 }
255 }
256
257 self.delete_position_buffer = Some(delete_position_buffer);
258 self.pk_index_state_table.try_flush().await?;
259 }
260
261 async fn apply_resolver_chunk(&mut self, chunk: StreamChunk) -> StreamExecutorResult<()> {
262 for (op, row) in chunk.rows() {
263 if op != Op::Insert {
264 bail!(
265 "iceberg pk-index writer {} expected resolver inserts, got {op:?}",
266 self.sink_id
267 );
268 }
269 self.pk_index_state_table.insert(row);
270 }
271 self.pk_index_state_table.try_flush().await?;
272 Ok(())
273 }
274
275 #[try_stream(ok = Message, error = StreamExecutorError)]
276 async fn checkpoint_barrier(&mut self, barrier: Barrier) {
277 barrier.assume_no_update_vnode_bitmap(self.ctx.id)?;
278 let mut metadata = None;
279 if barrier.is_checkpoint() {
280 if let Some(chunk) = self
281 .delete_position_buffer
282 .take()
283 .and_then(|mut builder| builder.consume_all())
284 {
285 yield Message::Chunk(chunk.into());
286 }
287 metadata = self.writer.flush().await?;
288 }
289
290 self.pk_index_state_table
291 .commit_assert_no_update_vnode_bitmap(barrier.epoch)
292 .await?;
293 if let Some(metadata) = metadata
294 && metadata.metadata.is_some()
295 {
296 self.local_barrier_manager
297 .report_iceberg_pk_index_sink_metadata(
298 barrier.epoch,
299 self.sink_id,
300 self.ctx.id,
301 PbIcebergPkIndexSinkRole::Writer,
302 Some(metadata),
303 );
304 }
305 yield Message::Barrier(barrier);
306 }
307
308 fn validate_compaction_barrier(
309 &self,
310 barrier: &Barrier,
311 expected_task: IcebergCompactionTaskId,
312 expected_phase: Phase,
313 expected_prev: u64,
314 ) -> StreamExecutorResult<()> {
315 if !barrier.is_checkpoint() || barrier.epoch.prev != expected_prev {
316 bail!(
317 "iceberg pk-index writer {} expected checkpoint {:?} starting at {}, got {:?}",
318 self.sink_id,
319 expected_phase,
320 expected_prev,
321 barrier
322 );
323 }
324 match barrier.iceberg_pk_index_compaction() {
325 Some(context)
326 if context.sink_id == self.sink_id
327 && context.task_id == expected_task
328 && context.phase == expected_phase as i32 =>
329 {
330 Ok(())
331 }
332 _ => bail!(
333 "iceberg pk-index writer {} expected matching {:?} context for task {}, got {:?}",
334 self.sink_id,
335 expected_phase,
336 expected_task,
337 barrier
338 ),
339 }
340 }
341
342 fn validate_aligned_barriers(
343 &self,
344 left: &Barrier,
345 right: &Barrier,
346 ) -> StreamExecutorResult<()> {
347 let context_matches = match (
348 left.iceberg_pk_index_compaction(),
349 right.iceberg_pk_index_compaction(),
350 ) {
351 (None, None) => true,
352 (Some(left), Some(right)) => {
353 left.sink_id == right.sink_id
354 && left.task_id == right.task_id
355 && left.phase == right.phase
356 }
357 _ => false,
358 };
359 if left.epoch != right.epoch
360 || left.kind != right.kind
361 || left.mutation != right.mutation
362 || !context_matches
363 {
364 bail!(
365 "iceberg pk-index writer {} received mismatched left/right barriers: left={:?}, right={:?}",
366 self.sink_id,
367 left,
368 right
369 );
370 }
371 Ok(())
372 }
373
374 fn compaction_begin(
375 &self,
376 barrier: &Barrier,
377 ) -> StreamExecutorResult<Option<IcebergCompactionTaskId>> {
378 let context = match barrier.iceberg_pk_index_compaction() {
379 Some(context) if context.sink_id == self.sink_id => context,
380 _ => return Ok(None),
381 };
382 if context.phase == Phase::End as i32 {
383 bail!(
384 "iceberg pk-index writer {} received unexpected End in Normal mode for task {}",
385 self.sink_id,
386 context.task_id
387 );
388 }
389 if context.phase != Phase::Begin as i32 {
390 bail!(
391 "iceberg pk-index writer {} expected Begin context for task {}, got {:?}",
392 self.sink_id,
393 context.task_id,
394 context.phase
395 );
396 }
397 Ok(Some(context.task_id))
398 }
399
400 #[try_stream(ok = Message, error = StreamExecutorError)]
401 async fn execute_inner(mut self) {
402 let mut input = self.input.take().unwrap().execute();
403 let mut resolver_input = self.resolver_input.take().unwrap().execute();
404
405 let barrier = expect_first_barrier(&mut input).await?;
407 let remap_first = expect_first_barrier(&mut resolver_input).await?;
408 self.validate_aligned_barriers(&barrier, &remap_first)?;
409 let first_epoch = barrier.epoch;
410 yield Message::Barrier(barrier);
411 self.pk_index_state_table.init_epoch(first_epoch).await?;
412
413 loop {
414 let mode = std::mem::replace(&mut self.mode, WriterInputMode::Normal);
415 match mode {
416 WriterInputMode::Normal => {
417 let mut completed = false;
418 #[for_await]
419 for msg in self.execute_normal(&mut input, &mut resolver_input, &mut completed)
420 {
421 yield msg?;
422 }
423 if completed {
424 break;
425 }
426 }
427 WriterInputMode::ResolvingRight {
428 task_id,
429 begin_epoch,
430 } => {
431 #[for_await]
432 for msg in
433 self.execute_resolving_right(&mut resolver_input, task_id, begin_epoch)
434 {
435 yield msg?;
436 }
437 }
438 WriterInputMode::DrainingLeft { task_id, barrier } => {
439 #[for_await]
440 for msg in self.execute_draining_left(&mut input, task_id, barrier) {
441 yield msg?;
442 }
443 }
444 }
445 }
446 }
447
448 #[try_stream(ok = Message, error = StreamExecutorError)]
449 async fn execute_normal<'a>(
450 &'a mut self,
451 input: &'a mut BoxedMessageStream,
452 resolver_input: &'a mut BoxedMessageStream,
453 completed: &'a mut bool,
454 ) {
455 #[for_await]
456 for msg in input {
457 match msg? {
458 Message::Chunk(chunk) =>
459 {
460 #[for_await]
461 for chunk in self.process_chunk(chunk) {
462 yield Message::Chunk(chunk?.into());
463 }
464 }
465 Message::Watermark(watermark) => {
466 yield Message::Watermark(watermark);
467 }
468 Message::Barrier(barrier) => {
469 let msg = next_msg(resolver_input).await?;
470 let remap_barrier = match msg {
471 Message::Barrier(b) => b,
472 _ => bail!(
473 "iceberg pk-index writer {} expected barrier on remap input, got {msg:?}",
474 self.sink_id
475 ),
476 };
477 self.validate_aligned_barriers(&barrier, &remap_barrier)?;
478 let begin = self.compaction_begin(&barrier)?;
479 let epoch = barrier.epoch;
480 #[for_await]
481 for msg in self.checkpoint_barrier(barrier) {
482 yield msg?;
483 }
484
485 if let Some(task_id) = begin {
486 self.mode = WriterInputMode::ResolvingRight {
487 task_id,
488 begin_epoch: epoch,
489 };
490 return Ok(());
491 }
492 }
493 }
494 }
495
496 *completed = true;
497 }
498
499 #[try_stream(ok = Message, error = StreamExecutorError)]
500 async fn execute_resolving_right<'a>(
501 &'a mut self,
502 resolver_input: &'a mut BoxedMessageStream,
503 task_id: IcebergCompactionTaskId,
504 begin_epoch: EpochPair,
505 ) {
506 #[for_await]
507 for msg in resolver_input {
508 match msg? {
509 Message::Chunk(chunk) => {
510 self.apply_resolver_chunk(chunk).await?;
511 }
512 Message::Watermark(_) => bail!(
513 "iceberg pk-index writer {} received watermark on remap input while resolving task {}",
514 self.sink_id,
515 task_id
516 ),
517 Message::Barrier(barrier) => {
518 self.validate_compaction_barrier(
519 &barrier,
520 task_id,
521 Phase::End,
522 begin_epoch.curr,
523 )?;
524 self.mode = WriterInputMode::DrainingLeft { task_id, barrier };
525 return Ok(());
526 }
527 }
528 }
529
530 bail!(
531 "iceberg pk-index writer {} right input closed before End for task {}",
532 self.sink_id,
533 task_id
534 );
535 }
536
537 #[try_stream(ok = Message, error = StreamExecutorError)]
538 async fn execute_draining_left<'a>(
539 &'a mut self,
540 input: &'a mut BoxedMessageStream,
541 task_id: IcebergCompactionTaskId,
542 expected: Barrier,
543 ) {
544 #[for_await]
545 for msg in input {
546 match msg? {
547 Message::Chunk(chunk) =>
548 {
549 #[for_await]
550 for chunk in self.process_chunk(chunk) {
551 yield Message::Chunk(chunk?.into());
552 }
553 }
554 Message::Watermark(_) => bail!(
555 "iceberg pk-index writer {} received watermark on left input while draining task {}",
556 self.sink_id,
557 task_id
558 ),
559 Message::Barrier(barrier) => {
560 self.validate_aligned_barriers(&barrier, &expected)?;
561 #[for_await]
562 for msg in self.checkpoint_barrier(barrier) {
563 yield msg?;
564 }
565 self.mode = WriterInputMode::Normal;
566 return Ok(());
567 }
568 }
569 }
570
571 bail!(
572 "iceberg pk-index writer {} left input closed before End for task {}",
573 self.sink_id,
574 task_id
575 );
576 }
577}
578
579impl<S, W> Execute for WriterExecutor<S, W>
580where
581 S: StateStore,
582 W: IcebergWriter,
583{
584 fn execute(self: Box<Self>) -> BoxedMessageStream {
585 self.execute_inner().boxed()
586 }
587}
588
589async fn next_msg(input: &mut BoxedMessageStream) -> StreamExecutorResult<Message> {
590 input
591 .next()
592 .await
593 .ok_or_else(|| {
594 StreamExecutorError::channel_closed(
595 "iceberg pk-index writer input channel closed unexpectedly",
596 )
597 })
598 .flatten()
599}
600
601#[cfg(test)]
602#[path = "writer_test.rs"]
603mod tests;