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::{
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
51/// Leaf executor for the iceberg pk-index coordinated compaction.
52///
53/// It has **no stream input**. One resolver runs in a transient independent job and is driven by
54/// barriers delivered directly through `barrier_receiver`. It only produces resolved SURVIVOR
55/// output chunks that a downstream `Writer` consumes to rebuild the index.
56///
57/// Lifecycle:
58/// - On the first barrier: forward it (no state to init).
59/// - After resolve completes, on the next barrier (the end barrier): REPORT the conflict DV
60///   files to meta, forward the barrier, and exit.
61///
62/// # Output chunk contract
63///
64/// Schema: `[pk_columns.., file_path: Varchar, position: Int64]`. Every emitted row is
65/// a SURVIVOR: `Writer` upserts `pk -> (file_path, position)` in the index (index-only; the data
66/// already lives in the compaction output file). Conflicts (rows killed during the window) are NOT
67/// emitted here.
68pub 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    /// Data types of the survivor output chunk: `[pk_columns.., file_path, position]`.
158    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        // Deletion-vector positions per input data file, at R (baseline) and at N (post-window).
194        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        // Selectively re-read the input data files at the delete-diff positions to recover the primary
198        // keys deleted during the window. `pk` datums are keyed by an `OwnedRow` so membership tests
199        // against the output scan below are exact.
200        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        // 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(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        // First barrier: forward it. There is no state table to initialize.
264        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
311/// Walks the manifests of `snapshot` and collects, per input data file, the merged set of deleted
312/// positions from every live position-delete file referencing it.
313async 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        // Only delete manifests can carry position-delete files.
328        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
357/// Builds a projection for PK root columns and maps the projected physical order back to the
358/// downstream PK order.
359fn 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    // A fresh uuid per resolve keeps file names unique across the vnode-aligned resolver actors.
497    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        // Compaction outputs are not committed yet and the resolver only receives their paths.
513        // Meta retains the full output data-file metadata through B2 and backfills each
514        // resolver/merger delete file's partition from its referenced data file during pre-commit.
515        // Therefore these provisional delete artifacts intentionally omit `PartitionKey` here.
516        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        // The output files are uncommitted, so no pre-existing DV file is being overwritten.
547        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}