1use anyhow::{Context, anyhow};
16use hashbrown::{HashMap, HashSet};
17use iceberg::delete_vector::DeleteVector;
18use iceberg::io::FileIO;
19use iceberg::spec::{DataContentType, ManifestContentType, SnapshotRef};
20use iceberg::table::Table;
21use iceberg::writer::file_writer::location_generator::DefaultLocationGenerator;
22use parquet::arrow::{ParquetRecordBatchStreamBuilder, ProjectionMask};
23use parquet::schema::types::SchemaDescriptor;
24use risingwave_common::array::DataChunk;
25use risingwave_common::array::arrow::IcebergArrowConvert;
26use risingwave_common::id::SinkId;
27use risingwave_common::row::RowExt;
28use risingwave_common::util::chunk_coalesce::DataChunkBuilder;
29use risingwave_connector::sink::SinkError;
30use risingwave_connector::sink::iceberg::{
31 IcebergConfig, IcebergPositionDeleteCommitResult, PositionDeleteFileNameGenerators,
32 read_position_deletes_from_file, serialize_data_files_default_spec, write_position_delete_file,
33};
34use risingwave_connector::source::iceberg::ParquetFileReader;
35use risingwave_pb::connector_service::SinkMetadata;
36use risingwave_pb::id::IcebergCompactionTaskId;
37use risingwave_pb::stream_plan::iceberg_pk_index_compaction_context::{Phase, ResolverTaskInput};
38use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
39use risingwave_rpc_client::MetaClient;
40use tokio::sync::mpsc::UnboundedReceiver;
41use uuid::Uuid;
42
43use super::load_table_at_least;
44use crate::executor::prelude::*;
45use crate::task::LocalBarrierManager;
46
47pub struct CompactionResolverExecutor {
65 ctx: ActorContextRef,
66 sink_id: SinkId,
67 iceberg_config: IcebergConfig,
68 pk_indices: Vec<usize>,
69 pk_data_types: Vec<DataType>,
70 chunk_size: usize,
71 local_barrier_manager: LocalBarrierManager,
72 barrier_receiver: UnboundedReceiver<Barrier>,
73 meta_client: MetaClient,
74}
75
76fn resolver_task_from_initial_barrier(
77 sink_id: SinkId,
78 barrier: &Barrier,
79) -> StreamExecutorResult<(IcebergCompactionTaskId, ResolverTaskInput)> {
80 match barrier.iceberg_pk_index_compaction() {
81 Some(context)
82 if barrier.is_checkpoint()
83 && context.sink_id == sink_id
84 && context.phase == Phase::Begin as i32 =>
85 {
86 let task_input = context.resolver_task_input.clone().ok_or_else(|| {
87 StreamExecutorError::from(anyhow!(
88 "compaction resolver sink {} task {} missing resolver task input",
89 sink_id,
90 context.task_id
91 ))
92 })?;
93 Ok((context.task_id, task_input))
94 }
95 _ => Err(StreamExecutorError::from(anyhow!(
96 "compaction resolver sink {} expected initial begin barrier, got {:?}",
97 sink_id,
98 barrier
99 ))),
100 }
101}
102
103fn validate_resolver_end_barrier(
104 sink_id: SinkId,
105 barrier: &Barrier,
106 begin: &Barrier,
107 task_id: IcebergCompactionTaskId,
108) -> StreamExecutorResult<()> {
109 match barrier.iceberg_pk_index_compaction() {
110 Some(context)
111 if barrier.is_checkpoint()
112 && barrier.epoch.prev == begin.epoch.curr
113 && context.sink_id == sink_id
114 && context.task_id == task_id
115 && context.phase == Phase::End as i32
116 && context.resolver_task_input.is_none() =>
117 {
118 Ok(())
119 }
120 _ => Err(StreamExecutorError::from(anyhow!(
121 "compaction resolver sink {} task {} expected end barrier, got {:?}",
122 sink_id,
123 task_id,
124 barrier
125 ))),
126 }
127}
128
129impl CompactionResolverExecutor {
130 #[expect(clippy::too_many_arguments)]
131 pub fn new(
132 ctx: ActorContextRef,
133 sink_id: SinkId,
134 iceberg_config: IcebergConfig,
135 pk_indices: Vec<usize>,
136 pk_data_types: Vec<DataType>,
137 chunk_size: usize,
138 local_barrier_manager: LocalBarrierManager,
139 barrier_receiver: UnboundedReceiver<Barrier>,
140 meta_client: MetaClient,
141 ) -> Self {
142 Self {
143 ctx,
144 sink_id,
145 iceberg_config,
146 pk_indices,
147 pk_data_types,
148 chunk_size,
149 local_barrier_manager,
150 barrier_receiver,
151 meta_client,
152 }
153 }
154
155 fn output_data_types(&self) -> Vec<DataType> {
157 let mut types = self.pk_data_types.clone();
158 types.push(DataType::Varchar);
159 types.push(DataType::Int64);
160 types
161 }
162
163 #[try_stream(ok = DataChunk, error = SinkError)]
164 async fn resolve<'a>(
165 &'a self,
166 task_input: &'a ResolverTaskInput,
167 expected_snapshot: Option<i64>,
168 conflict_delete_metadata: &'a mut Option<SinkMetadata>,
169 ) {
170 let table = load_table_at_least(&self.iceberg_config, expected_snapshot).await?;
171 let file_io = table.file_io();
172 let snapshot_r = table
173 .metadata()
174 .snapshot_by_id(task_input.read_snapshot_id)
175 .with_context(|| {
176 format!(
177 "compaction read snapshot {} not present in table metadata",
178 task_input.read_snapshot_id
179 )
180 })?;
181 let snapshot_n = table
182 .metadata()
183 .current_snapshot()
184 .context("table has no current snapshot to resolve compaction against")?;
185
186 let input_paths: HashSet<&str> = task_input
187 .input_data_file_paths
188 .iter()
189 .map(String::as_str)
190 .collect();
191
192 let dv_r = collect_input_dvs(&table, snapshot_r, &input_paths).await?;
194 let dv_n = collect_input_dvs(&table, snapshot_n, &input_paths).await?;
195
196 let mut delete_diff_pk_set: HashSet<OwnedRow> = HashSet::new();
200 let stream = futures::stream::iter(task_input.input_data_file_paths.iter().map(
201 |input_path| async {
202 let Some(n_positions) = dv_n.get(input_path) else {
203 return Ok(Vec::new());
204 };
205 let r_positions = dv_r.get(input_path);
206 let mut diff = n_positions.clone();
207 if let Some(r_positions) = r_positions {
208 diff -= r_positions;
209 }
210 if diff.is_empty() {
211 return Ok(Vec::new());
212 }
213 scan_input_pks_at_positions(file_io, input_path, &self.pk_indices, &diff).await
214 },
215 ))
216 .buffer_unordered(8);
217 #[for_await]
218 for pks in stream {
219 delete_diff_pk_set.extend(pks?);
220 }
221
222 let mut conflicts: HashMap<String, DeleteVector> = HashMap::new();
224 #[for_await]
225 for chunk in self.scan_output_file(
226 file_io,
227 &task_input.output_data_file_paths,
228 &delete_diff_pk_set,
229 &mut conflicts,
230 ) {
231 yield chunk?;
232 }
233
234 *conflict_delete_metadata =
235 write_conflict_delete_files(&self.iceberg_config, &table, conflicts).await?;
236 }
237
238 #[try_stream(ok = DataChunk, error = SinkError)]
239 async fn scan_output_file<'a>(
240 &'a self,
241 file_io: &'a FileIO,
242 output_data_file_paths: &'a [String],
243 delete_diff_pk_set: &'a HashSet<OwnedRow>,
244 conflicts: &'a mut HashMap<String, DeleteVector>,
245 ) {
246 let mut buffer = DataChunkBuilder::new(self.output_data_types(), self.chunk_size);
247 for file in output_data_file_paths {
248 #[for_await]
249 for chunk in scan_output_file_inner(
250 file_io,
251 file,
252 &self.pk_indices,
253 delete_diff_pk_set,
254 &mut buffer,
255 conflicts,
256 ) {
257 yield chunk?;
258 }
259 }
260 if let Some(chunk) = buffer.consume_all() {
261 yield chunk;
262 }
263 }
264
265 #[try_stream(ok = Message, error = StreamExecutorError)]
266 async fn execute_inner(mut self) {
267 let actor_id = self.ctx.id;
268
269 let first_barrier = self.barrier_receiver.recv().await.ok_or_else(|| {
271 StreamExecutorError::channel_closed("compaction resolver barrier receiver")
272 })?;
273 let (task_id, task_input) =
274 resolver_task_from_initial_barrier(self.sink_id, &first_barrier)?;
275 let begin_barrier = first_barrier.clone();
276 yield Message::Barrier(first_barrier);
277
278 let expected_snapshot = self
279 .meta_client
280 .wait_iceberg_pk_index_sink_epoch(self.sink_id, begin_barrier.epoch.prev)
281 .await?;
282
283 let mut conflict_delete_metadata = None;
284 #[for_await]
285 for chunk in self.resolve(
286 &task_input,
287 expected_snapshot,
288 &mut conflict_delete_metadata,
289 ) {
290 let chunk = chunk.map_err(|e| (e, self.sink_id))?;
291 yield Message::Chunk(chunk.into());
292 }
293
294 let Some(barrier) = self.barrier_receiver.recv().await else {
295 return Err(StreamExecutorError::channel_closed(
296 "compaction resolver barrier receiver",
297 ));
298 };
299 validate_resolver_end_barrier(self.sink_id, &barrier, &begin_barrier, task_id)?;
300 if let Some(metadata) = conflict_delete_metadata
301 && metadata.metadata.is_some()
302 {
303 self.local_barrier_manager
304 .report_iceberg_pk_index_sink_metadata(
305 barrier.epoch,
306 self.sink_id,
307 actor_id,
308 PbIcebergPkIndexSinkRole::CompactionResolver,
309 Some(metadata),
310 );
311 }
312 yield Message::Barrier(barrier);
313 }
314}
315
316impl Execute for CompactionResolverExecutor {
317 fn execute(self: Box<Self>) -> BoxedMessageStream {
318 self.execute_inner().boxed()
319 }
320}
321
322async fn collect_input_dvs(
325 table: &Table,
326 snapshot: &SnapshotRef,
327 input_paths: &HashSet<&str>,
328) -> Result<HashMap<String, DeleteVector>, SinkError> {
329 let file_io = table.file_io();
330 let mut map: HashMap<String, DeleteVector> = HashMap::new();
331
332 let manifest_list = table
333 .object_cache()
334 .get_manifest_list(snapshot, &table.metadata_ref())
335 .await?;
336
337 for manifest_file in manifest_list.entries() {
338 if manifest_file.content != ManifestContentType::Deletes {
340 continue;
341 }
342 let manifest = manifest_file.load_manifest(file_io).await?;
343 for entry in manifest.entries() {
344 if !entry.is_alive() {
345 continue;
346 }
347 let data_file = entry.data_file();
348 if data_file.content_type() == DataContentType::PositionDeletes
349 && let Some(referenced) = data_file.referenced_data_file()
350 && input_paths.contains(referenced.as_str())
351 {
352 let positions = read_position_deletes_from_file(file_io, data_file)
353 .await
354 .map_err(SinkError::Iceberg)?;
355 if map.insert(referenced, positions).is_some() {
356 return Err(SinkError::Iceberg(anyhow!(
357 "input data file {} has multiple live position-delete files in snapshot {}",
358 data_file.referenced_data_file().unwrap(),
359 snapshot.snapshot_id()
360 )));
361 }
362 }
363 }
364 }
365 Ok(map)
366}
367
368fn pk_projection(
371 parquet_schema: &SchemaDescriptor,
372 pk_indices: &[usize],
373) -> Result<(ProjectionMask, Vec<usize>), SinkError> {
374 let root_count = parquet_schema.root_schema().get_fields().len();
375 let mut physical_indices = Vec::with_capacity(pk_indices.len());
376 let mut seen = HashSet::with_capacity(pk_indices.len());
377
378 for &pk_index in pk_indices {
379 if pk_index >= root_count {
380 return Err(SinkError::Iceberg(anyhow!(
381 "pk index {pk_index} is out of range for parquet schema with {root_count} root columns"
382 )));
383 }
384 if !seen.insert(pk_index) {
385 return Err(SinkError::Iceberg(anyhow!("duplicate pk index {pk_index}")));
386 }
387 physical_indices.push(pk_index);
388 }
389
390 physical_indices.sort_unstable();
391 let pk_order = pk_indices
392 .iter()
393 .map(|pk_index| {
394 physical_indices
395 .binary_search(pk_index)
396 .expect("validated pk index must be projected")
397 })
398 .collect();
399 let projection = ProjectionMask::roots(parquet_schema, physical_indices);
400
401 Ok((projection, pk_order))
402}
403
404async fn scan_input_pks_at_positions(
405 file_io: &FileIO,
406 path: &str,
407 pk_indices: &[usize],
408 want_positions: &DeleteVector,
409) -> Result<Vec<OwnedRow>, SinkError> {
410 let mut results = Vec::new();
411 let input_file = file_io.new_input(path)?;
412 let metadata = input_file.metadata().await?;
413 let reader = input_file.reader().await?;
414 let builder = ParquetRecordBatchStreamBuilder::new(ParquetFileReader::new(metadata, reader))
415 .await
416 .map_err(|e| anyhow!(e).context(format!("open parquet reader for {path}")))?;
417 let (projection, pk_order) = pk_projection(builder.parquet_schema(), pk_indices)?;
418 let mut stream = builder
419 .with_projection(projection)
420 .build()
421 .map_err(|e| anyhow!(e).context(format!("build parquet stream for {path}")))?;
422
423 let mut base_pos = 0;
424 let mut iter = want_positions.iter().peekable();
425 while let Some(batch) = stream.next().await {
426 let batch =
427 batch.map_err(|e| anyhow!(e).context(format!("read parquet batch of {path}")))?;
428 let chunk = IcebergArrowConvert
429 .chunk_from_record_batch(&batch)
430 .map_err(|e| anyhow!(e).context(format!("convert parquet batch of {path}")))?
431 .project(&pk_order);
432 while let Some(&iter_pos) = iter.peek() {
433 let chunk_pos = iter_pos as usize - base_pos;
434 if chunk_pos >= chunk.capacity() {
435 break;
436 }
437 let pk = chunk.row_at(chunk_pos).0.to_owned_row();
438 results.push(pk);
439 iter.next();
440 }
441 base_pos += chunk.capacity();
442 }
443
444 Ok(results)
445}
446
447#[try_stream(ok = DataChunk, error = SinkError)]
448async fn scan_output_file_inner<'a>(
449 file_io: &'a FileIO,
450 path: &'a str,
451 pk_indices: &'a [usize],
452 delete_diff_pk_set: &'a HashSet<OwnedRow>,
453 buffer: &'a mut DataChunkBuilder,
454 conflicts: &'a mut HashMap<String, DeleteVector>,
455) {
456 let input_file = file_io.new_input(path)?;
457 let metadata = input_file.metadata().await?;
458 let reader = input_file.reader().await?;
459 let builder = ParquetRecordBatchStreamBuilder::new(ParquetFileReader::new(metadata, reader))
460 .await
461 .map_err(|e| anyhow!(e).context(format!("open parquet reader for {path}")))?;
462 let (projection, pk_order) = pk_projection(builder.parquet_schema(), pk_indices)?;
463 let stream = builder
464 .with_projection(projection)
465 .build()
466 .map_err(|e| anyhow!(e).context(format!("build parquet stream for {path}")))?;
467
468 let mut position: i64 = 0;
469
470 #[for_await]
471 for batch in stream {
472 let batch =
473 batch.map_err(|e| anyhow!(e).context(format!("read parquet batch of {path}")))?;
474 let chunk = IcebergArrowConvert
475 .chunk_from_record_batch(&batch)
476 .map_err(|e| anyhow!(e).context(format!("convert parquet batch of {path}")))?
477 .project(&pk_order);
478 for row_idx in 0..chunk.capacity() {
479 let pk = chunk.row_at(row_idx).0.into_owned_row();
480 let row_position = position + row_idx as i64;
481 if delete_diff_pk_set.contains(&pk) {
482 conflicts
483 .entry_ref(path)
484 .or_default()
485 .insert(row_position as u64);
486 } else if let Some(full_chunk) = buffer.append_one_row(pk.chain([
487 Some(ScalarRefImpl::Utf8(path)),
488 Some(ScalarRefImpl::Int64(row_position)),
489 ])) {
490 yield full_chunk;
491 }
492 }
493 position += chunk.capacity() as i64;
494 }
495}
496
497async fn write_conflict_delete_files(
498 config: &IcebergConfig,
499 table: &Table,
500 conflicts: HashMap<String, DeleteVector>,
501) -> Result<Option<SinkMetadata>, SinkError> {
502 if conflicts.is_empty() {
503 return Ok(None);
504 }
505
506 let location_generator = DefaultLocationGenerator::new(table.metadata())?;
507 let uuid_suffix = Uuid::now_v7();
508 let file_name_generators = PositionDeleteFileNameGenerators::new(uuid_suffix);
509 let format_version = table.metadata().format_version();
510
511 let mut delete_files = Vec::with_capacity(conflicts.len());
512 for (output_file, delete_vector) in conflicts {
513 let file = write_position_delete_file(
519 table,
520 config,
521 &location_generator,
522 &file_name_generators,
523 format_version,
524 output_file,
525 &delete_vector,
526 None,
527 )
528 .await?;
529 delete_files.push(file);
530 }
531
532 let delete_files = serialize_data_files_default_spec(table, delete_files)?;
533 let result = IcebergPositionDeleteCommitResult {
534 schema_id: table.metadata().current_schema_id(),
535 partition_spec_id: table.metadata().default_partition_spec_id(),
536 delete_files,
537 overwrite_files: vec![],
539 };
540 let sink_metadata = SinkMetadata::try_from(&result)?;
541 Ok(Some(sink_metadata))
542}
543
544#[cfg(test)]
545mod tests {
546 use std::sync::Arc;
547
548 use parquet::basic::Type as PhysicalType;
549 use parquet::schema::types::{SchemaDescriptor, Type};
550 use risingwave_common::array::DataChunkTestExt;
551 use risingwave_common::util::epoch::test_epoch;
552 use risingwave_pb::stream_plan::IcebergPkIndexCompactionContext;
553
554 use super::*;
555
556 fn test_parquet_schema() -> SchemaDescriptor {
557 SchemaDescriptor::new(Arc::new(
558 Type::group_type_builder("schema")
559 .with_fields(
560 ["column_0", "column_1", "column_2", "column_3", "column_4"]
561 .into_iter()
562 .map(|name| {
563 Arc::new(
564 Type::primitive_type_builder(name, PhysicalType::INT32)
565 .build()
566 .unwrap(),
567 )
568 })
569 .collect(),
570 )
571 .build()
572 .unwrap(),
573 ))
574 }
575
576 #[test]
577 fn test_pk_projection_restores_downstream_pk_order() {
578 let parquet_schema = test_parquet_schema();
579
580 let (_, order) = pk_projection(&parquet_schema, &[2, 0, 1]).unwrap();
581 assert_eq!(order, vec![2, 0, 1]);
582
583 let (_, order) = pk_projection(&parquet_schema, &[4, 1, 3]).unwrap();
584 assert_eq!(order, vec![2, 0, 1]);
585 }
586
587 #[test]
588 fn test_pk_projection_reorders_heterogeneous_chunk() {
589 let chunk = DataChunk::from_pretty(
590 "i T
591 7 k",
592 );
593 let projected = chunk.project(&[1, 0]);
594
595 assert_eq!(
596 projected,
597 DataChunk::from_pretty(
598 "T i
599 k 7",
600 )
601 );
602 }
603
604 fn compaction_barrier(
605 epoch: u64,
606 phase: Phase,
607 resolver_task_input: Option<ResolverTaskInput>,
608 ) -> Barrier {
609 Barrier::new_test_barrier(test_epoch(epoch)).with_iceberg_pk_index_compaction(
610 IcebergPkIndexCompactionContext {
611 sink_id: SinkId::new(7),
612 task_id: 11.into(),
613 phase: phase as i32,
614 resolver_task_input,
615 },
616 )
617 }
618
619 #[test]
620 fn test_resolver_task_is_initialized_from_b1() {
621 let task_input = ResolverTaskInput {
622 output_data_file_paths: vec!["output.parquet".into()],
623 input_data_file_paths: vec!["input.parquet".into()],
624 read_snapshot_id: 42,
625 };
626 let barrier = compaction_barrier(2, Phase::Begin, Some(task_input.clone()));
627
628 let (task_id, actual) =
629 resolver_task_from_initial_barrier(SinkId::new(7), &barrier).unwrap();
630 assert_eq!(task_id, IcebergCompactionTaskId::from(11));
631 assert_eq!(actual, task_input);
632 }
633
634 #[test]
635 fn test_resolver_requires_task_input_only_on_b1() {
636 let begin = compaction_barrier(2, Phase::Begin, Some(ResolverTaskInput::default()));
637 let end = compaction_barrier(3, Phase::End, None);
638 validate_resolver_end_barrier(SinkId::new(7), &end, &begin, 11.into()).unwrap();
639
640 let invalid_begin = compaction_barrier(2, Phase::Begin, None);
641 assert!(resolver_task_from_initial_barrier(SinkId::new(7), &invalid_begin).is_err());
642 let invalid_end = compaction_barrier(3, Phase::End, Some(ResolverTaskInput::default()));
643 assert!(
644 validate_resolver_end_barrier(SinkId::new(7), &invalid_end, &begin, 11.into()).is_err()
645 );
646 }
647}