risingwave_meta/manager/iceberg_pk_index_sink/
mod.rs1mod 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#[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
125pub 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
139pub 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}