Skip to main content

risingwave_meta/manager/iceberg_pk_index_sink/
mod.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
15//! Manager for the Iceberg pk-index sink path. Owns per-sink commit coordinators that
16//! drive iceberg `commit_epoch` ahead of hummock `commit_epoch` and persist
17//! exactly-once state via `pending_sink_state`.
18//!
19//! This is intentionally separate from [`crate::manager::sink_coordination`]
20//! (which serves V1/V2 sinks via gRPC). Future responsibilities such as
21//! per-sink compaction will live alongside the per-sink commit coordinator here.
22
23mod committed_epoch;
24mod coordinator;
25mod manager;
26
27use std::collections::{BTreeMap, HashMap};
28
29use anyhow::anyhow;
30use iceberg::spec::SerializedDataFile;
31pub use manager::IcebergPkIndexSinkManager;
32use risingwave_common::secret::LocalSecretManager;
33use risingwave_connector::sink::catalog::SinkId;
34use risingwave_connector::sink::iceberg::{ENABLE_PK_INDEX, IcebergConfig};
35use risingwave_connector::source::UPSTREAM_SOURCE_KEY;
36use risingwave_pb::catalog::PbSink;
37use risingwave_pb::stream_service::barrier_complete_response::IcebergPkIndexSinkMetadata as PbIcebergPkIndexSinkMetadata;
38
39#[derive(educe::Educe)]
40#[educe(Debug)]
41pub(crate) struct CompactionOverwrite {
42    pub sink_id: SinkId,
43    pub epoch: u64,
44    pub schema_id: i32,
45    pub partition_spec_id: i32,
46    #[educe(Debug(ignore))]
47    pub output_files: Vec<SerializedDataFile>,
48    pub input_file_paths: Vec<String>,
49    pub read_snapshot_id: i64,
50}
51
52/// Metadata collected for one sink/epoch before the Hummock checkpoint is committed.
53/// Ordinary reports and the optional compaction overwrite share one transport shape so the
54/// barrier path does not need separate metadata variants.
55#[derive(Debug)]
56pub(crate) struct IcebergPkIndexPreCommitMetadata {
57    pub sink_id: SinkId,
58    pub prev_epoch: u64,
59    pub reports: Vec<PbIcebergPkIndexSinkMetadata>,
60    pub compaction: Option<CompactionOverwrite>,
61}
62
63impl From<PbIcebergPkIndexSinkMetadata> for IcebergPkIndexPreCommitMetadata {
64    fn from(report: PbIcebergPkIndexSinkMetadata) -> Self {
65        Self {
66            sink_id: report.sink_id,
67            prev_epoch: report.prev_epoch,
68            reports: vec![report],
69            compaction: None,
70        }
71    }
72}
73
74impl From<CompactionOverwrite> for IcebergPkIndexPreCommitMetadata {
75    fn from(overwrite: CompactionOverwrite) -> Self {
76        Self {
77            sink_id: overwrite.sink_id,
78            prev_epoch: overwrite.epoch,
79            reports: vec![],
80            compaction: Some(overwrite),
81        }
82    }
83}
84
85pub(crate) fn group_pre_commit_metadata(
86    metadata: Vec<IcebergPkIndexPreCommitMetadata>,
87) -> anyhow::Result<Vec<IcebergPkIndexPreCommitMetadata>> {
88    let mut grouped = HashMap::<SinkId, IcebergPkIndexPreCommitMetadata>::new();
89    for metadata in metadata {
90        let IcebergPkIndexPreCommitMetadata {
91            sink_id,
92            prev_epoch,
93            reports,
94            compaction,
95        } = metadata;
96        let input = grouped
97            .entry(sink_id)
98            .or_insert_with(|| IcebergPkIndexPreCommitMetadata {
99                sink_id,
100                prev_epoch,
101                reports: Vec::new(),
102                compaction: None,
103            });
104        if input.prev_epoch != prev_epoch {
105            anyhow::bail!(
106                "iceberg v3 sink {} pre-commit metadata disagrees on prev_epoch: {} vs {}",
107                sink_id,
108                input.prev_epoch,
109                prev_epoch
110            );
111        }
112        input.reports.extend(reports);
113        if let Some(overwrite) = compaction
114            && input.compaction.replace(overwrite).is_some()
115        {
116            anyhow::bail!(
117                "iceberg v3 sink {} has multiple compaction overwrites in one barrier",
118                sink_id
119            );
120        }
121    }
122    Ok(grouped.into_values().collect())
123}
124
125/// Returns true if the given sink properties identify a Iceberg pk-index sink
126/// (i.e. an iceberg sink with `enable_pk_index = 'true'`).
127pub fn is_iceberg_pk_index_sink(properties: &BTreeMap<String, String>) -> bool {
128    let connector_match = properties
129        .get(UPSTREAM_SOURCE_KEY)
130        .map(|v| v.eq_ignore_ascii_case("iceberg"))
131        .unwrap_or(false);
132    let pk_index_enabled = properties
133        .get(ENABLE_PK_INDEX)
134        .map(|v| v.eq_ignore_ascii_case("true"))
135        .unwrap_or(false);
136    connector_match && pk_index_enabled
137}
138
139/// Build an [`IcebergConfig`] from a [`PbSink`], filling secret refs along the
140/// way. Used at CREATE SINK time and during recovery to (re-)register the
141/// commit coordinator.
142pub fn build_iceberg_config(pb_sink: &PbSink) -> anyhow::Result<IcebergConfig> {
143    let properties: BTreeMap<String, String> = pb_sink.properties.clone().into_iter().collect();
144    let secret_refs: BTreeMap<_, _> = pb_sink.secret_refs.clone().into_iter().collect();
145    let with_secrets = LocalSecretManager::global()
146        .fill_secrets(properties, secret_refs)
147        .map_err(|e| anyhow!(e).context("fill secrets for iceberg"))?;
148    IcebergConfig::from_btreemap(with_secrets)
149        .map_err(|e| anyhow!(e).context("parse iceberg config"))
150}
151
152#[cfg(test)]
153mod tests {
154    use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
155
156    use super::*;
157
158    fn report(
159        sink_id: u32,
160        prev_epoch: u64,
161        role: PbIcebergPkIndexSinkRole,
162    ) -> PbIcebergPkIndexSinkMetadata {
163        PbIcebergPkIndexSinkMetadata {
164            sink_id: SinkId::new(sink_id),
165            prev_epoch,
166            role: role as i32,
167            ..Default::default()
168        }
169    }
170
171    fn overwrite(sink_id: u32, epoch: u64) -> CompactionOverwrite {
172        CompactionOverwrite {
173            sink_id: SinkId::new(sink_id),
174            epoch,
175            schema_id: 3,
176            partition_spec_id: 4,
177            output_files: vec![],
178            input_file_paths: vec![],
179            read_snapshot_id: 1,
180        }
181    }
182
183    #[test]
184    fn group_pre_commit_metadata_combines_reports_and_compaction() {
185        let inputs = group_pre_commit_metadata(vec![
186            report(7, 60, PbIcebergPkIndexSinkRole::Writer).into(),
187            overwrite(7, 60).into(),
188            report(7, 60, PbIcebergPkIndexSinkRole::PositionDeleteMerger).into(),
189        ])
190        .unwrap();
191
192        assert_eq!(inputs.len(), 1);
193        let input = &inputs[0];
194        assert_eq!(input.sink_id, SinkId::new(7));
195        assert_eq!(input.prev_epoch, 60);
196        assert_eq!(input.reports.len(), 2);
197        let overwrite = input.compaction.as_ref().unwrap();
198        assert_eq!(overwrite.schema_id, 3);
199        assert_eq!(overwrite.partition_spec_id, 4);
200    }
201
202    #[test]
203    fn group_pre_commit_metadata_preserves_ordinary_only_sink() {
204        let inputs =
205            group_pre_commit_metadata(vec![report(8, 61, PbIcebergPkIndexSinkRole::Writer).into()])
206                .unwrap();
207
208        assert_eq!(inputs.len(), 1);
209        assert_eq!(inputs[0].sink_id, SinkId::new(8));
210        assert_eq!(inputs[0].prev_epoch, 61);
211        assert_eq!(inputs[0].reports.len(), 1);
212        assert!(inputs[0].compaction.is_none());
213    }
214
215    #[test]
216    fn group_pre_commit_metadata_rejects_duplicate_compaction() {
217        let error =
218            group_pre_commit_metadata(vec![overwrite(7, 60).into(), overwrite(7, 60).into()])
219                .unwrap_err();
220        assert!(error.to_string().contains("multiple compaction overwrites"));
221    }
222
223    #[test]
224    fn group_pre_commit_metadata_rejects_epoch_mismatch() {
225        let error = group_pre_commit_metadata(vec![
226            report(7, 59, PbIcebergPkIndexSinkRole::Writer).into(),
227            overwrite(7, 60).into(),
228        ])
229        .unwrap_err();
230        assert!(error.to_string().contains("disagrees on prev_epoch"));
231    }
232}