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