Skip to main content

risingwave_stream/executor/iceberg_with_pk_index/
compaction_resolver.rs

1// Copyright 2026 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
47/// Leaf executor for the iceberg pk-index coordinated compaction.
48///
49/// It has **no stream input**. The actor is part of the sink's static graph but only runs while a
50/// compaction is being applied. It is driven by barriers delivered directly through
51/// `barrier_receiver` and produces resolved survivor rows for the downstream writer.
52///
53/// Lifecycle:
54/// - On the first barrier: forward it (no state to init).
55/// - After resolve completes, on the next barrier (the end barrier): REPORT the conflict DV
56///   files to meta, forward the barrier, and exit.
57///
58/// # Output chunk contract
59///
60/// Schema: `[pk_columns.., file_path: Varchar, position: Int64]`. Every emitted row is
61/// a SURVIVOR: `Writer` upserts `pk -> (file_path, position)` in the index (index-only; the data
62/// already lives in the compaction output file). Conflicts (rows killed during the window) are NOT
63/// emitted here.
64pub 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    /// Data types of the survivor output chunk: `[pk_columns.., file_path, position]`.
156    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        // Deletion-vector positions per input data file, at R (baseline) and at N (post-window).
193        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        // Selectively re-read the input data files at the delete-diff positions to recover the primary
197        // keys deleted during the window. `pk` datums are keyed by an `OwnedRow` so membership tests
198        // against the output scan below are exact.
199        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        // Scan every output data file once per actor.
223        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        // First barrier: forward it. There is no state table to initialize.
270        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
322/// Walks the manifests of `snapshot` and collects, per input data file, the merged set of deleted
323/// positions from every live position-delete file referencing it.
324async 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        // Only delete manifests can carry position-delete files.
339        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
368/// Builds a projection for PK root columns and maps the projected physical order back to the
369/// downstream PK order.
370fn 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        // Compaction outputs are not committed yet and the resolver only receives their paths.
514        // Meta retains the full output data-file metadata through the resolver detach checkpoint
515        // and backfills each resolver/merger delete file's partition from its referenced data file
516        // during pre-commit.
517        // Therefore these provisional delete artifacts intentionally omit `PartitionKey` here.
518        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        // The output files are uncommitted, so no pre-existing DV file is being overwritten.
538        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}