1use std::collections::HashMap;
16use std::sync::Arc;
17use std::time::Duration;
18
19use anyhow::{Context, anyhow};
20use async_trait::async_trait;
21use iceberg::Catalog;
22use iceberg::arrow::schema_to_arrow_schema;
23use iceberg::spec::{DataFile, FormatVersion, Operation, SerializedDataFile, TableMetadata};
24use iceberg::table::Table;
25use iceberg::transaction::{AddColumn, ApplyTransactionAction, FastAppendAction, Transaction};
26use itertools::Itertools;
27use risingwave_common::array::arrow::arrow_schema_iceberg::{
28 DataType as ArrowDataType, Field as ArrowField, Fields as ArrowFields,
29};
30use risingwave_common::array::arrow::{IcebergArrowConvert, IcebergCreateTableArrowConvert};
31use risingwave_common::bail;
32use risingwave_common::catalog::{Field, RISINGWAVE_ICEBERG_COMMIT_EPOCH};
33use risingwave_common::error::IcebergError;
34use risingwave_pb::connector_service::SinkMetadata;
35use risingwave_pb::connector_service::sink_metadata::Metadata::Serialized;
36use risingwave_pb::connector_service::sink_metadata::SerializedMetadata;
37use risingwave_pb::stream_plan::PbSinkSchemaChange;
38use serde::{Deserialize, Serialize};
39use serde_json::from_value;
40use thiserror_ext::AsReport;
41use tokio::sync::mpsc::UnboundedSender;
42use tracing::warn;
43
44use super::commit_retry::{self, CommitError, CommitRetryLogContext};
45use super::{GLOBAL_SINK_METRICS, IcebergConfig, SinkError, commit_branch, resolve_partition_type};
46use crate::connector_common::{IcebergCommittedSnapshot, IcebergSinkCompactionUpdate};
47use crate::sink::catalog::SinkId;
48use crate::sink::{Result, SinglePhaseCommitCoordinator, SinkParam, TwoPhaseCommitCoordinator};
49
50const SCHEMA_ID: &str = "schema_id";
51const PARTITION_SPEC_ID: &str = "partition_spec_id";
52const DATA_FILES: &str = "data_files";
53
54#[derive(Default, Clone)]
55pub struct IcebergCommitResult {
56 pub schema_id: i32,
57 pub partition_spec_id: i32,
58 pub data_files: Vec<SerializedDataFile>,
59}
60
61impl IcebergCommitResult {
62 pub fn try_from(value: &SinkMetadata) -> Result<Self> {
63 let Some(Serialized(value)) = &value.metadata else {
64 bail!("Can't create iceberg sink write result from empty data!");
65 };
66
67 Self::try_from_serialized_bytes(&value.metadata)
68 }
69
70 pub fn try_from_serialized_bytes(value: &[u8]) -> Result<Self> {
71 let mut values = if let serde_json::Value::Object(value) =
72 serde_json::from_slice::<serde_json::Value>(value)
73 .context("Can't parse iceberg sink metadata")?
74 {
75 value
76 } else {
77 bail!("iceberg sink metadata should be an object");
78 };
79
80 let schema_id;
81 if let Some(serde_json::Value::Number(value)) = values.remove(SCHEMA_ID) {
82 schema_id = value
83 .as_u64()
84 .ok_or_else(|| anyhow!("schema_id should be a u64"))?;
85 } else {
86 bail!("iceberg sink metadata should have schema_id");
87 }
88
89 let partition_spec_id;
90 if let Some(serde_json::Value::Number(value)) = values.remove(PARTITION_SPEC_ID) {
91 partition_spec_id = value
92 .as_u64()
93 .ok_or_else(|| anyhow!("partition_spec_id should be a u64"))?;
94 } else {
95 bail!("iceberg sink metadata should have partition_spec_id");
96 }
97
98 let data_files: Vec<SerializedDataFile>;
99 if let serde_json::Value::Array(values) = values
100 .remove(DATA_FILES)
101 .ok_or_else(|| anyhow!("iceberg sink metadata should have data_files object"))?
102 {
103 data_files = values
104 .into_iter()
105 .map(from_value::<SerializedDataFile>)
106 .collect::<std::result::Result<_, _>>()
107 .unwrap();
108 } else {
109 bail!("iceberg sink metadata should have data_files object");
110 }
111
112 Ok(Self {
113 schema_id: schema_id as i32,
114 partition_spec_id: partition_spec_id as i32,
115 data_files,
116 })
117 }
118}
119
120impl<'a> TryFrom<&'a IcebergCommitResult> for SinkMetadata {
121 type Error = SinkError;
122
123 fn try_from(value: &'a IcebergCommitResult) -> std::result::Result<SinkMetadata, Self::Error> {
124 let bytes = <Vec<u8>>::try_from(value)?;
125 Ok(SinkMetadata {
126 metadata: Some(Serialized(SerializedMetadata { metadata: bytes })),
127 })
128 }
129}
130
131impl<'a> TryFrom<&'a IcebergCommitResult> for Vec<u8> {
132 type Error = SinkError;
133
134 fn try_from(value: &'a IcebergCommitResult) -> std::result::Result<Vec<u8>, Self::Error> {
135 let json_data_files = serde_json::Value::Array(
136 value
137 .data_files
138 .iter()
139 .map(serde_json::to_value)
140 .collect::<std::result::Result<Vec<serde_json::Value>, _>>()
141 .context("Can't serialize data files to json")?,
142 );
143 let json_value = serde_json::Value::Object(
144 vec![
145 (
146 SCHEMA_ID.to_owned(),
147 serde_json::Value::Number(value.schema_id.into()),
148 ),
149 (
150 PARTITION_SPEC_ID.to_owned(),
151 serde_json::Value::Number(value.partition_spec_id.into()),
152 ),
153 (DATA_FILES.to_owned(), json_data_files),
154 ]
155 .into_iter()
156 .collect(),
157 );
158 Ok(serde_json::to_vec(&json_value).context("Can't serialize iceberg sink metadata")?)
159 }
160}
161
162#[derive(Default, Clone, Serialize, Deserialize)]
163pub struct IcebergPositionDeleteCommitResult {
164 pub schema_id: i32,
165 pub partition_spec_id: i32,
166 pub delete_files: Vec<SerializedDataFile>,
167 pub overwrite_files: Vec<SerializedDataFile>,
168}
169
170impl<'a> TryFrom<&'a SinkMetadata> for IcebergPositionDeleteCommitResult {
171 type Error = SinkError;
172
173 fn try_from(value: &'a SinkMetadata) -> Result<Self> {
174 let Some(Serialized(value)) = &value.metadata else {
175 bail!("Can't create iceberg dv merger commit result from empty data!");
176 };
177 let value = serde_json::from_slice(&value.metadata)
178 .context("Can't deserialize iceberg dv merger commit result from metadata")?;
179 Ok(value)
180 }
181}
182
183impl<'a> TryFrom<&'a IcebergPositionDeleteCommitResult> for SinkMetadata {
184 type Error = SinkError;
185
186 fn try_from(value: &'a IcebergPositionDeleteCommitResult) -> Result<SinkMetadata> {
187 let bytes = serde_json::to_vec(value)
188 .context("Can't serialize iceberg dv merger commit result to metadata")?;
189 Ok(SinkMetadata {
190 metadata: Some(Serialized(SerializedMetadata { metadata: bytes })),
191 })
192 }
193}
194
195fn arrow_data_type_compatible(current: &ArrowDataType, expected: &ArrowDataType) -> bool {
196 use ArrowDataType::*;
197
198 match (current, expected) {
199 (Decimal128(_, _), Decimal128(_, _)) => true,
205 (Binary, LargeBinary) | (LargeBinary, Binary) => true,
206 (List(current_field), List(expected_field)) => {
207 arrow_data_type_compatible(current_field.data_type(), expected_field.data_type())
208 }
209 (Map(current_field, current_sorted), Map(expected_field, expected_sorted)) => {
210 current_sorted == expected_sorted
211 && arrow_data_type_compatible(current_field.data_type(), expected_field.data_type())
212 }
213 (Struct(current_fields), Struct(expected_fields)) => {
214 let expected_fields = expected_fields
215 .iter()
216 .map(|field| field.as_ref().clone())
217 .collect_vec();
218 schema_contains_same_fields(current_fields, &expected_fields)
219 }
220 _ => current == expected,
221 }
222}
223
224fn schema_contains_same_fields(current: &ArrowFields, expected: &[ArrowField]) -> bool {
225 if current.len() != expected.len() {
226 return false;
227 }
228
229 let mut unmatched_current = current.iter().collect_vec();
230 expected.iter().all(|expected_field| {
231 let Some(pos) = unmatched_current.iter().position(|current_field| {
232 current_field.name() == expected_field.name()
233 && arrow_data_type_compatible(current_field.data_type(), expected_field.data_type())
234 }) else {
235 return false;
236 };
237 unmatched_current.swap_remove(pos);
238 true
239 })
240}
241
242pub struct IcebergSinkCommitter {
243 pub(super) catalog: Arc<dyn Catalog>,
244 pub(super) table: Table,
245 pub last_commit_epoch: u64,
246 pub(crate) sink_id: SinkId,
247 pub(crate) config: IcebergConfig,
248 pub(crate) param: SinkParam,
249 pub(super) commit_retry_num: u32,
250 pub(crate) iceberg_compact_stat_sender: Option<UnboundedSender<IcebergSinkCompactionUpdate>>,
251}
252
253impl IcebergSinkCommitter {
254 fn latest_observed_snapshot(&self) -> Option<IcebergCommittedSnapshot> {
255 let branch = commit_branch(self.config.r#type.as_str(), self.config.write_mode);
256 let metadata = self.table.metadata();
257 metadata
258 .snapshot_for_ref(&branch)
259 .map(|snapshot| IcebergCommittedSnapshot {
260 branch,
261 snapshot_id: snapshot.snapshot_id(),
262 timestamp_ms: snapshot.timestamp_ms(),
263 max_file_sequence_number: (metadata.format_version() >= FormatVersion::V2)
264 .then_some(snapshot.sequence_number()),
265 })
266 }
267
268 fn notify_iceberg_compaction_scheduler(&self, force_compaction: bool) {
269 let Some(iceberg_compact_stat_sender) = &self.iceberg_compact_stat_sender else {
270 return;
271 };
272 let branch = commit_branch(self.config.r#type.as_str(), self.config.write_mode);
273
274 let Some(observed_snapshot) = self.latest_observed_snapshot() else {
275 warn!(
276 iceberg_component = "sink_committer",
277 iceberg_operation = "notify_compaction",
278 sink_id = %self.sink_id,
279 sink_name = %self.param.sink_name,
280 table = %self.table.identifier(),
281 branch = %branch,
282 force_compaction,
283 "iceberg_sink_compaction_update_skipped",
284 );
285 return;
286 };
287
288 let observed_snapshot_id = observed_snapshot.snapshot_id;
289 let observed_snapshot_timestamp_ms = observed_snapshot.timestamp_ms;
290 let observed_snapshot_branch = observed_snapshot.branch.clone();
291
292 if iceberg_compact_stat_sender
293 .send(IcebergSinkCompactionUpdate {
294 sink_id: self.sink_id,
295 force_compaction,
296 observed_snapshot,
297 })
298 .is_err()
299 {
300 warn!(
301 iceberg_component = "sink_committer",
302 iceberg_operation = "notify_compaction",
303 sink_id = %self.sink_id,
304 sink_name = %self.param.sink_name,
305 table = %self.table.identifier(),
306 force_compaction,
307 observed_snapshot_id,
308 observed_snapshot_timestamp_ms,
309 observed_snapshot_branch = %observed_snapshot_branch,
310 "iceberg_sink_compaction_update_send_failed",
311 );
312 }
313 }
314}
315
316#[async_trait]
317impl SinglePhaseCommitCoordinator for IcebergSinkCommitter {
318 async fn init(&mut self) -> Result<()> {
319 tracing::info!(
320 iceberg_component = "sink_committer",
321 iceberg_operation = "init",
322 sink_id = %self.param.sink_id,
323 sink_name = %self.param.sink_name,
324 table = %self.table.identifier(),
325 "iceberg_sink_committer_initialized",
326 );
327
328 Ok(())
329 }
330
331 async fn commit_data(&mut self, epoch: u64, metadata: Vec<SinkMetadata>) -> Result<()> {
332 tracing::debug!(
333 iceberg_component = "sink_committer",
334 iceberg_operation = "direct_commit",
335 sink_id = %self.sink_id,
336 sink_name = %self.param.sink_name,
337 table = %self.table.identifier(),
338 epoch,
339 metadata_count = metadata.len(),
340 "iceberg_sink_commit_started",
341 );
342
343 if metadata.is_empty() {
344 tracing::debug!(
345 iceberg_component = "sink_committer",
346 iceberg_operation = "direct_commit",
347 sink_id = %self.sink_id,
348 table = %self.table.identifier(),
349 epoch,
350 "iceberg_sink_commit_skipped_empty_metadata",
351 );
352 return Ok(());
353 }
354
355 if let Some((write_results, snapshot_id)) = self.pre_commit_inner(epoch, metadata)? {
357 self.commit_data_impl(epoch, write_results, snapshot_id)
358 .await?;
359 }
360
361 Ok(())
362 }
363
364 async fn commit_schema_change(
365 &mut self,
366 epoch: u64,
367 schema_change: PbSinkSchemaChange,
368 ) -> Result<()> {
369 tracing::info!(
370 iceberg_component = "sink_committer",
371 iceberg_operation = "schema_change",
372 sink_id = %self.sink_id,
373 sink_name = %self.param.sink_name,
374 table = %self.table.identifier(),
375 epoch,
376 schema_change = ?schema_change,
377 "iceberg_sink_schema_change_commit_started",
378 );
379 self.commit_schema_change_impl(schema_change).await?;
380 tracing::info!(
381 iceberg_component = "sink_committer",
382 iceberg_operation = "schema_change",
383 sink_id = %self.sink_id,
384 table = %self.table.identifier(),
385 epoch,
386 "iceberg_sink_schema_change_commit_succeeded",
387 );
388
389 Ok(())
390 }
391}
392
393#[async_trait]
394impl TwoPhaseCommitCoordinator for IcebergSinkCommitter {
395 async fn init(&mut self) -> Result<()> {
396 tracing::info!(
397 iceberg_component = "sink_committer",
398 iceberg_operation = "init",
399 sink_id = %self.param.sink_id,
400 sink_name = %self.param.sink_name,
401 table = %self.table.identifier(),
402 "iceberg_sink_committer_initialized",
403 );
404
405 Ok(())
406 }
407
408 async fn pre_commit(
409 &mut self,
410 epoch: u64,
411 metadata: Vec<SinkMetadata>,
412 _schema_change: Option<PbSinkSchemaChange>,
413 ) -> Result<Option<Vec<u8>>> {
414 tracing::debug!(
415 iceberg_component = "sink_committer",
416 iceberg_operation = "pre_commit",
417 sink_id = %self.sink_id,
418 sink_name = %self.param.sink_name,
419 table = %self.table.identifier(),
420 epoch,
421 metadata_count = metadata.len(),
422 "iceberg_sink_pre_commit_started",
423 );
424
425 let (write_results, snapshot_id) = match self.pre_commit_inner(epoch, metadata)? {
426 Some((write_results, snapshot_id)) => (write_results, snapshot_id),
427 None => {
428 tracing::debug!(
429 iceberg_component = "sink_committer",
430 iceberg_operation = "pre_commit",
431 sink_id = %self.sink_id,
432 table = %self.table.identifier(),
433 epoch,
434 "iceberg_sink_pre_commit_skipped_no_data",
435 );
436 return Ok(None);
437 }
438 };
439
440 let mut write_results_bytes = Vec::new();
441 for each_parallelism_write_result in write_results {
442 let each_parallelism_write_result_bytes =
443 <Vec<u8>>::try_from(&each_parallelism_write_result)?;
444 write_results_bytes.push(each_parallelism_write_result_bytes);
445 }
446
447 let snapshot_id_bytes: Vec<u8> = snapshot_id.to_le_bytes().to_vec();
448 write_results_bytes.push(snapshot_id_bytes);
449
450 let pre_commit_metadata_bytes: Vec<u8> = serialize_metadata(write_results_bytes);
451 tracing::debug!(
452 iceberg_component = "sink_committer",
453 iceberg_operation = "pre_commit",
454 sink_id = %self.sink_id,
455 table = %self.table.identifier(),
456 epoch,
457 snapshot_id,
458 pre_commit_metadata_bytes = pre_commit_metadata_bytes.len(),
459 "iceberg_sink_pre_commit_metadata_encoded",
460 );
461 Ok(Some(pre_commit_metadata_bytes))
462 }
463
464 async fn commit_data(&mut self, epoch: u64, commit_metadata: Vec<u8>) -> Result<()> {
465 tracing::debug!(
466 iceberg_component = "sink_committer",
467 iceberg_operation = "commit",
468 sink_id = %self.sink_id,
469 sink_name = %self.param.sink_name,
470 table = %self.table.identifier(),
471 epoch,
472 commit_metadata_bytes = commit_metadata.len(),
473 "iceberg_sink_commit_started",
474 );
475
476 if commit_metadata.is_empty() {
477 tracing::debug!(
478 iceberg_component = "sink_committer",
479 iceberg_operation = "commit",
480 sink_id = %self.sink_id,
481 table = %self.table.identifier(),
482 epoch,
483 "iceberg_sink_commit_skipped_empty_metadata",
484 );
485 return Ok(());
486 }
487
488 let mut payload = deserialize_metadata(commit_metadata);
490 if payload.is_empty() {
491 return Err(SinkError::Iceberg(anyhow!(
492 "Invalid commit metadata: empty payload"
493 )));
494 }
495
496 let snapshot_id_bytes = payload.pop().ok_or_else(|| {
498 SinkError::Iceberg(anyhow!("Invalid commit metadata: missing snapshot_id"))
499 })?;
500 let snapshot_id = i64::from_le_bytes(
501 snapshot_id_bytes
502 .try_into()
503 .map_err(|_| SinkError::Iceberg(anyhow!("Invalid snapshot id bytes")))?,
504 );
505
506 let write_results = payload
508 .into_iter()
509 .map(|p| IcebergCommitResult::try_from_serialized_bytes(&p))
510 .collect::<Result<Vec<_>>>()?;
511
512 let snapshot_committed = self.is_snapshot_id_in_iceberg(snapshot_id).await?;
513
514 if snapshot_committed {
515 tracing::info!(
516 iceberg_component = "sink_committer",
517 iceberg_operation = "commit",
518 sink_id = %self.sink_id,
519 sink_name = %self.param.sink_name,
520 table = %self.table.identifier(),
521 epoch,
522 snapshot_id,
523 "iceberg_sink_commit_skipped_snapshot_already_committed",
524 );
525 return Ok(());
526 }
527
528 self.commit_data_impl(epoch, write_results, snapshot_id)
529 .await
530 }
531
532 async fn commit_schema_change(
533 &mut self,
534 epoch: u64,
535 schema_change: PbSinkSchemaChange,
536 ) -> Result<()> {
537 let schema_updated = self.check_schema_change_applied(&schema_change)?;
538 if schema_updated {
539 tracing::info!(
540 iceberg_component = "sink_committer",
541 iceberg_operation = "schema_change",
542 sink_id = %self.sink_id,
543 table = %self.table.identifier(),
544 epoch,
545 "iceberg_sink_schema_change_skipped_already_applied",
546 );
547 return Ok(());
548 }
549
550 tracing::info!(
551 iceberg_component = "sink_committer",
552 iceberg_operation = "schema_change",
553 sink_id = %self.sink_id,
554 sink_name = %self.param.sink_name,
555 table = %self.table.identifier(),
556 epoch,
557 schema_change = ?schema_change,
558 "iceberg_sink_schema_change_commit_started",
559 );
560 self.commit_schema_change_impl(schema_change).await?;
561 tracing::info!(
562 iceberg_component = "sink_committer",
563 iceberg_operation = "schema_change",
564 sink_id = %self.sink_id,
565 table = %self.table.identifier(),
566 epoch,
567 "iceberg_sink_schema_change_commit_succeeded",
568 );
569
570 Ok(())
571 }
572
573 async fn abort(&mut self, epoch: u64, _commit_metadata: Vec<u8>) {
574 tracing::debug!(
576 iceberg_component = "sink_committer",
577 iceberg_operation = "abort",
578 sink_id = %self.sink_id,
579 table = %self.table.identifier(),
580 epoch,
581 "iceberg_sink_commit_abort_unimplemented",
582 );
583 }
584}
585
586impl IcebergSinkCommitter {
588 fn pre_commit_inner(
589 &mut self,
590 epoch: u64,
591 metadata: Vec<SinkMetadata>,
592 ) -> Result<Option<(Vec<IcebergCommitResult>, i64)>> {
593 let write_results: Vec<IcebergCommitResult> = metadata
594 .iter()
595 .map(IcebergCommitResult::try_from)
596 .collect::<Result<Vec<IcebergCommitResult>>>()?;
597 let data_file_count: usize = write_results.iter().map(|r| r.data_files.len()).sum();
598 tracing::debug!(
599 iceberg_component = "sink_committer",
600 iceberg_operation = "pre_commit",
601 sink_id = %self.sink_id,
602 table = %self.table.identifier(),
603 epoch,
604 writer_result_count = write_results.len(),
605 data_file_count,
606 "iceberg_sink_pre_commit_metadata_decoded",
607 );
608
609 if write_results.is_empty() || write_results.iter().all(|r| r.data_files.is_empty()) {
611 return Ok(None);
612 }
613
614 let expect_schema_id = write_results[0].schema_id;
615 let expect_partition_spec_id = write_results[0].partition_spec_id;
616
617 if write_results
619 .iter()
620 .any(|r| r.schema_id != expect_schema_id)
621 || write_results
622 .iter()
623 .any(|r| r.partition_spec_id != expect_partition_spec_id)
624 {
625 return Err(SinkError::Iceberg(anyhow!(
626 "schema_id and partition_spec_id should be the same in all write results"
627 )));
628 }
629
630 let snapshot_id = FastAppendAction::generate_snapshot_id(&self.table);
631 tracing::debug!(
632 iceberg_component = "sink_committer",
633 iceberg_operation = "pre_commit",
634 sink_id = %self.sink_id,
635 table = %self.table.identifier(),
636 epoch,
637 snapshot_id,
638 schema_id = expect_schema_id,
639 partition_spec_id = expect_partition_spec_id,
640 data_file_count,
641 "iceberg_sink_pre_commit_snapshot_assigned",
642 );
643
644 Ok(Some((write_results, snapshot_id)))
645 }
646
647 async fn commit_data_impl(
648 &mut self,
649 epoch: u64,
650 write_results: Vec<IcebergCommitResult>,
651 snapshot_id: i64,
652 ) -> Result<()> {
653 assert!(
655 !write_results.is_empty() && !write_results.iter().all(|r| r.data_files.is_empty())
656 );
657
658 self.wait_for_snapshot_limit().await?;
660
661 let expect_schema_id = write_results[0].schema_id;
662 let expect_partition_spec_id = write_results[0].partition_spec_id;
663 let data_file_count: usize = write_results.iter().map(|r| r.data_files.len()).sum();
664 let target_branch = commit_branch(self.config.r#type.as_str(), self.config.write_mode);
665 tracing::info!(
666 iceberg_component = "sink_committer",
667 iceberg_operation = "commit",
668 sink_id = %self.sink_id,
669 sink_name = %self.param.sink_name,
670 table = %self.table.identifier(),
671 branch = %target_branch,
672 epoch,
673 snapshot_id,
674 schema_id = expect_schema_id,
675 partition_spec_id = expect_partition_spec_id,
676 data_file_count,
677 retry_num = self.commit_retry_num,
678 "iceberg_sink_commit_applying",
679 );
680
681 self.table = commit_retry::reload_table(
683 self.catalog.as_ref(),
684 self.table.identifier(),
685 expect_schema_id,
686 expect_partition_spec_id,
687 )
688 .await
689 .map_err(SinkError::Iceberg)?;
690
691 let Some(schema) = self.table.metadata().schema_by_id(expect_schema_id) else {
692 return Err(SinkError::Iceberg(anyhow!(
693 "Can't find schema by id {}",
694 expect_schema_id
695 )));
696 };
697 let partition_type = resolve_partition_type(&self.table, expect_partition_spec_id, schema)?;
698
699 let data_files = write_results
700 .into_iter()
701 .flat_map(|r| {
702 r.data_files.into_iter().map(|f| {
703 f.try_into(expect_partition_spec_id, &partition_type, schema)
704 .map_err(|err| SinkError::Iceberg(anyhow!(err)))
705 })
706 })
707 .collect::<Result<Vec<DataFile>>>()?;
708
709 let catalog = self.catalog.clone();
714 let sink_id = self.sink_id;
715 let table_ident = self.table.identifier().clone();
716 let table_name = table_ident.to_string();
717 let retry_log_context = CommitRetryLogContext::new(
718 "sink_committer",
719 "commit",
720 table_name.clone(),
721 target_branch.clone(),
722 )
723 .with_sink_id(sink_id)
724 .with_epoch(epoch)
725 .with_snapshot_id(snapshot_id);
726
727 let table = commit_retry::run_with_retry(
728 catalog.clone(),
729 table_ident.clone(),
730 expect_schema_id,
731 expect_partition_spec_id,
732 self.commit_retry_num as usize,
733 retry_log_context,
734 |table| {
735 let table_name = table_name.clone();
736 let target_branch = target_branch.clone();
737 let data_files = data_files.clone();
738 let catalog = catalog.clone();
739 async move {
740 let txn = Transaction::new(&table);
741 let append_action = txn
742 .fast_append()
743 .set_snapshot_id(snapshot_id)
744 .set_target_branch(target_branch.clone())
745 .set_snapshot_properties(HashMap::from([(
746 RISINGWAVE_ICEBERG_COMMIT_EPOCH.to_owned(),
747 epoch.to_string(),
748 )]))
749 .add_data_files(data_files);
750
751 let tx = append_action.apply(txn).map_err(|err| {
752 let err: IcebergError = err.into();
753 tracing::error!(
754 iceberg_component = "sink_committer",
755 iceberg_operation = "commit",
756 sink_id = %sink_id,
757 table = %table_name,
758 epoch,
759 snapshot_id,
760 branch = %target_branch,
761 data_file_count,
762 error = %err.as_report(),
763 "iceberg_sink_commit_fast_append_apply_failed",
764 );
765 CommitError::Commit(anyhow!(err).context("apply iceberg fast_append"))
766 })?;
767
768 let table = tx.commit(catalog.as_ref()).await.map_err(|err| {
769 let err: IcebergError = err.into();
770 tracing::error!(
771 iceberg_component = "sink_committer",
772 iceberg_operation = "commit",
773 sink_id = %sink_id,
774 table = %table_name,
775 epoch,
776 snapshot_id,
777 branch = %target_branch,
778 data_file_count,
779 error = %err.as_report(),
780 "iceberg_sink_commit_transaction_failed",
781 );
782 CommitError::Commit(anyhow!(err).context("commit iceberg transaction"))
783 })?;
784 Ok(table)
785 }
786 },
787 )
788 .await
789 .map_err(SinkError::Iceberg)?;
790 self.table = table;
791
792 let snapshot_num = self.table.metadata().snapshots().count();
793 let catalog_name = self.config.common.catalog_name();
794 let table_name = self.table.identifier().to_string();
795 let metrics_labels = [&self.param.sink_name, &catalog_name, &table_name];
796 GLOBAL_SINK_METRICS
797 .iceberg_snapshot_num
798 .with_guarded_label_values(&metrics_labels)
799 .set(snapshot_num as i64);
800
801 tracing::debug!(
802 iceberg_component = "sink_committer",
803 iceberg_operation = "commit",
804 sink_id = %self.sink_id,
805 sink_name = %self.param.sink_name,
806 table = %self.table.identifier(),
807 branch = %target_branch,
808 epoch,
809 snapshot_id,
810 snapshot_num,
811 data_file_count,
812 "iceberg_sink_commit_succeeded",
813 );
814
815 self.notify_iceberg_compaction_scheduler(false);
816
817 Ok(())
818 }
819
820 async fn is_snapshot_id_in_iceberg(&self, snapshot_id: i64) -> Result<bool> {
824 let table = self
825 .catalog
826 .load_table(self.table.identifier())
827 .await
828 .map_err(|err| SinkError::Iceberg(anyhow!(err).context("reload iceberg table")))?;
829 if table.metadata().snapshot_by_id(snapshot_id).is_some() {
830 Ok(true)
831 } else {
832 Ok(false)
833 }
834 }
835
836 fn check_schema_change_applied(&self, schema_change: &PbSinkSchemaChange) -> Result<bool> {
839 let current_schema = self.table.metadata().current_schema();
840 let current_arrow_schema = schema_to_arrow_schema(current_schema.as_ref())
841 .context("Failed to convert schema")
842 .map_err(SinkError::Iceberg)?;
843
844 let iceberg_arrow_convert = IcebergArrowConvert;
845
846 let schema_matches = |expected: &[ArrowField]| {
847 schema_contains_same_fields(current_arrow_schema.fields(), expected)
848 };
849
850 let original_arrow_fields: Vec<ArrowField> = schema_change
851 .original_schema
852 .iter()
853 .map(|pb_field| {
854 let field = Field::from(pb_field);
855 iceberg_arrow_convert
856 .to_arrow_field(&field.name, &field.data_type)
857 .context("Failed to convert field to arrow")
858 .map_err(SinkError::Iceberg)
859 })
860 .collect::<Result<_>>()?;
861
862 if schema_matches(&original_arrow_fields) {
864 tracing::debug!(
865 "Current iceberg schema matches original_schema ({} columns); schema change not applied",
866 original_arrow_fields.len()
867 );
868 return Ok(false);
869 }
870
871 let expected_after_change = match schema_change.op.as_ref() {
872 Some(risingwave_pb::stream_plan::sink_schema_change::Op::AddColumns(
873 add_columns_op,
874 )) => {
875 let add_arrow_fields: Vec<ArrowField> = add_columns_op
876 .fields
877 .iter()
878 .map(|pb_field| {
879 let field = Field::from(pb_field);
880 iceberg_arrow_convert
881 .to_arrow_field(&field.name, &field.data_type)
882 .context("Failed to convert field to arrow")
883 .map_err(SinkError::Iceberg)
884 })
885 .collect::<Result<_>>()?;
886
887 let mut expected_after_change = original_arrow_fields;
888 expected_after_change.extend(add_arrow_fields);
889 expected_after_change
890 }
891 Some(risingwave_pb::stream_plan::sink_schema_change::Op::DropColumns(
892 drop_columns_op,
893 )) => original_arrow_fields
894 .into_iter()
895 .filter(|field| {
896 !drop_columns_op
897 .column_names
898 .iter()
899 .any(|name| name == field.name())
900 })
901 .collect_vec(),
902 _ => {
903 return Err(SinkError::Iceberg(anyhow!(
904 "Unsupported sink schema change op in iceberg sink: {:?}",
905 schema_change.op
906 )));
907 }
908 };
909
910 if schema_matches(&expected_after_change) {
912 tracing::debug!(
913 "Current iceberg schema matches changed schema ({} columns); schema change already applied",
914 expected_after_change.len()
915 );
916 return Ok(true);
917 }
918
919 Err(SinkError::Iceberg(anyhow!(
920 "Current iceberg schema does not match either original_schema ({} cols) or changed schema; cannot determine whether schema change is applied",
921 schema_change.original_schema.len()
922 )))
923 }
924
925 async fn commit_schema_change_impl(&mut self, schema_change: PbSinkSchemaChange) -> Result<()> {
929 let iceberg_create_table_arrow_convert = IcebergCreateTableArrowConvert::default();
931 let mut new_fields = Vec::new();
932
933 let mut drop_column_names = Vec::new();
934 match schema_change.op.as_ref() {
935 Some(risingwave_pb::stream_plan::sink_schema_change::Op::AddColumns(
936 add_columns_op,
937 )) => {
938 let add_columns = add_columns_op.fields.iter().map(Field::from).collect_vec();
939 for field in &add_columns {
940 let arrow_field = iceberg_create_table_arrow_convert
942 .to_arrow_field(&field.name, &field.data_type)
943 .with_context(|| {
944 format!("Failed to convert field '{}' to arrow", field.name)
945 })
946 .map_err(SinkError::Iceberg)?;
947
948 let iceberg_type = iceberg::arrow::arrow_type_to_type(arrow_field.data_type())
950 .map_err(|err| {
951 SinkError::Iceberg(
952 anyhow!(err)
953 .context("Failed to convert Arrow type to Iceberg type"),
954 )
955 })?;
956
957 new_fields.push(AddColumn::optional(&field.name, iceberg_type));
958 tracing::info!("Prepared field '{}' for schema change", field.name);
959 }
960 }
961 Some(risingwave_pb::stream_plan::sink_schema_change::Op::DropColumns(
962 drop_columns_op,
963 )) => {
964 drop_column_names = drop_columns_op.column_names.clone();
965 }
966 _ => {
967 return Err(SinkError::Iceberg(anyhow!(
968 "Unsupported sink schema change op in iceberg sink: {:?}",
969 schema_change.op
970 )));
971 }
972 }
973
974 tracing::info!(
976 "Committing schema change to catalog for table {}",
977 self.table.identifier()
978 );
979
980 let txn = Transaction::new(&self.table);
981 let action_fields_added = new_fields.len();
982 let mut action = txn.update_schema();
983 for field in new_fields {
984 action = action.add_column(field);
985 }
986 for column_name in &drop_column_names {
987 action = action.delete_column(column_name);
988 }
989
990 let updated_table = action
991 .apply(txn)
992 .context("Failed to apply schema update action")
993 .map_err(SinkError::Iceberg)?
994 .commit(self.catalog.as_ref())
995 .await
996 .context("Failed to commit table schema change")
997 .map_err(SinkError::Iceberg)?;
998
999 self.table = updated_table;
1000
1001 tracing::info!(
1002 "Successfully committed schema change, added {} columns and dropped {} columns from iceberg table",
1003 action_fields_added,
1004 drop_column_names.len()
1005 );
1006
1007 Ok(())
1008 }
1009
1010 fn count_snapshots_since_rewrite_in_metadata(metadata: &TableMetadata, branch: &str) -> usize {
1013 let mut snapshot_id = metadata
1015 .snapshot_for_ref(branch)
1016 .map(|snapshot| snapshot.snapshot_id());
1017 let mut count = 0;
1018
1019 while let Some(current_snapshot_id) = snapshot_id {
1021 let Some(snapshot) = metadata.snapshot_by_id(current_snapshot_id) else {
1022 break;
1023 };
1024
1025 if snapshot.summary().operation == Operation::Replace {
1027 break;
1029 }
1030
1031 count += 1;
1033 snapshot_id = snapshot.parent_snapshot_id();
1034 }
1035
1036 count
1037 }
1038
1039 fn count_snapshots_since_rewrite(&self) -> usize {
1041 let branch = commit_branch(self.config.r#type.as_str(), self.config.write_mode);
1042 Self::count_snapshots_since_rewrite_in_metadata(self.table.metadata(), branch.as_str())
1043 }
1044
1045 async fn wait_for_snapshot_limit(&mut self) -> Result<()> {
1047 if let Some(max_snapshots) = self.config.max_snapshots_num_before_compaction {
1048 loop {
1049 let current_count = self.count_snapshots_since_rewrite();
1050
1051 if current_count < max_snapshots {
1052 tracing::info!(
1053 "Snapshot count check passed: {} < {}",
1054 current_count,
1055 max_snapshots
1056 );
1057 break;
1058 }
1059
1060 tracing::info!(
1061 "Snapshot count {} exceeds limit {}, waiting...",
1062 current_count,
1063 max_snapshots
1064 );
1065
1066 self.notify_iceberg_compaction_scheduler(true);
1067
1068 tokio::time::sleep(Duration::from_secs(30)).await;
1070
1071 let table_ident = self.table.identifier().clone();
1073 self.table = self.catalog.load_table(&table_ident).await.map_err(|err| {
1074 SinkError::Iceberg(anyhow!(err).context("reload iceberg table"))
1075 })?;
1076 }
1077 }
1078 Ok(())
1079 }
1080}
1081
1082fn serialize_metadata(metadata: Vec<Vec<u8>>) -> Vec<u8> {
1083 serde_json::to_vec(&metadata).unwrap()
1084}
1085
1086fn deserialize_metadata(bytes: Vec<u8>) -> Vec<Vec<u8>> {
1087 serde_json::from_slice(&bytes).unwrap()
1088}
1089
1090#[cfg(test)]
1091mod tests {
1092 use std::collections::HashMap;
1093
1094 use iceberg::spec::{
1095 FormatVersion, MAIN_BRANCH, NestedField, PrimitiveType, Schema, Snapshot,
1096 SnapshotReference, SnapshotRetention, SortOrder, Summary, TableMetadataBuilder, Type,
1097 UnboundPartitionSpec,
1098 };
1099 use risingwave_common::array::arrow::arrow_schema_iceberg::{
1100 DataType as ArrowDataType, Field as ArrowField, FieldRef as ArrowFieldRef,
1101 Fields as ArrowFields, Schema as ArrowSchema,
1102 };
1103
1104 use super::*;
1105
1106 #[test]
1107 fn test_schema_contains_same_fields_allows_binary_large_binary() {
1108 let current_schema = ArrowSchema::new(vec![
1109 ArrowField::new("k", ArrowDataType::Int32, true),
1110 ArrowField::new("v", ArrowDataType::LargeBinary, true),
1111 ]);
1112 let expected_fields = vec![
1113 ArrowField::new("k", ArrowDataType::Int32, true),
1114 ArrowField::new("v", ArrowDataType::Binary, true),
1115 ];
1116
1117 assert!(schema_contains_same_fields(
1118 current_schema.fields(),
1119 &expected_fields
1120 ));
1121 }
1122
1123 #[test]
1124 fn test_schema_contains_same_fields_allows_nested_binary_large_binary() {
1125 let current_schema = ArrowSchema::new(vec![
1126 ArrowField::new(
1127 "s",
1128 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1129 "payload",
1130 ArrowDataType::LargeBinary,
1131 true,
1132 )])),
1133 true,
1134 ),
1135 ArrowField::new(
1136 "l",
1137 ArrowDataType::List(ArrowFieldRef::new(ArrowField::new_list_field(
1138 ArrowDataType::LargeBinary,
1139 true,
1140 ))),
1141 true,
1142 ),
1143 ArrowField::new_map(
1144 "m",
1145 "entries",
1146 ArrowFieldRef::new(ArrowField::new("key", ArrowDataType::Utf8, false)),
1147 ArrowFieldRef::new(ArrowField::new("value", ArrowDataType::LargeBinary, true)),
1148 false,
1149 true,
1150 ),
1151 ]);
1152 let expected_fields = vec![
1153 ArrowField::new(
1154 "s",
1155 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1156 "payload",
1157 ArrowDataType::Binary,
1158 true,
1159 )])),
1160 true,
1161 ),
1162 ArrowField::new(
1163 "l",
1164 ArrowDataType::List(ArrowFieldRef::new(ArrowField::new_list_field(
1165 ArrowDataType::Binary,
1166 true,
1167 ))),
1168 true,
1169 ),
1170 ArrowField::new_map(
1171 "m",
1172 "entries",
1173 ArrowFieldRef::new(ArrowField::new("key", ArrowDataType::Utf8, false)),
1174 ArrowFieldRef::new(ArrowField::new("value", ArrowDataType::Binary, true)),
1175 false,
1176 true,
1177 ),
1178 ];
1179
1180 assert!(schema_contains_same_fields(
1181 current_schema.fields(),
1182 &expected_fields
1183 ));
1184 }
1185
1186 #[test]
1187 fn test_schema_contains_same_fields_allows_decimal_precision_delta() {
1188 let current_schema = ArrowSchema::new(vec![ArrowField::new(
1189 "d",
1190 ArrowDataType::Decimal128(28, 10),
1191 true,
1192 )]);
1193 let expected_fields = vec![ArrowField::new(
1194 "d",
1195 ArrowDataType::Decimal128(38, 10),
1196 true,
1197 )];
1198
1199 assert!(schema_contains_same_fields(
1200 current_schema.fields(),
1201 &expected_fields
1202 ));
1203 }
1204
1205 #[test]
1206 fn test_schema_contains_same_fields_allows_decimal_precision_and_scale_delta() {
1207 let current_schema = ArrowSchema::new(vec![ArrowField::new(
1208 "d",
1209 ArrowDataType::Decimal128(38, 2),
1210 true,
1211 )]);
1212 let expected_fields = vec![ArrowField::new(
1213 "d",
1214 ArrowDataType::Decimal128(38, 10),
1215 true,
1216 )];
1217
1218 assert!(schema_contains_same_fields(
1219 current_schema.fields(),
1220 &expected_fields
1221 ));
1222 }
1223
1224 #[test]
1225 fn test_schema_contains_same_fields_rejects_length_mismatch() {
1226 let current_schema =
1227 ArrowSchema::new(vec![ArrowField::new("k", ArrowDataType::Int32, true)]);
1228 let expected_fields = vec![
1229 ArrowField::new("k", ArrowDataType::Int32, true),
1230 ArrowField::new("v", ArrowDataType::Utf8, true),
1231 ];
1232
1233 assert!(!schema_contains_same_fields(
1234 current_schema.fields(),
1235 &expected_fields
1236 ));
1237 }
1238
1239 #[test]
1240 fn test_schema_contains_same_fields_rejects_type_mismatch() {
1241 let current_schema =
1242 ArrowSchema::new(vec![ArrowField::new("v", ArrowDataType::Utf8, true)]);
1243 let expected_fields = vec![ArrowField::new("v", ArrowDataType::Int32, true)];
1244
1245 assert!(!schema_contains_same_fields(
1246 current_schema.fields(),
1247 &expected_fields
1248 ));
1249 }
1250
1251 #[test]
1252 fn test_schema_contains_same_fields_rejects_reused_duplicate_match() {
1253 let current_schema = ArrowSchema::new(vec![
1254 ArrowField::new("v", ArrowDataType::Int32, true),
1255 ArrowField::new("v", ArrowDataType::Utf8, true),
1256 ]);
1257 let expected_fields = vec![
1258 ArrowField::new("v", ArrowDataType::Int32, true),
1259 ArrowField::new("v", ArrowDataType::Int32, true),
1260 ];
1261
1262 assert!(!schema_contains_same_fields(
1263 current_schema.fields(),
1264 &expected_fields
1265 ));
1266 }
1267
1268 #[test]
1269 fn test_schema_contains_same_fields_rejects_nested_type_mismatch() {
1270 let current_schema = ArrowSchema::new(vec![ArrowField::new(
1271 "s",
1272 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1273 "payload",
1274 ArrowDataType::Utf8,
1275 true,
1276 )])),
1277 true,
1278 )]);
1279 let expected_fields = vec![ArrowField::new(
1280 "s",
1281 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1282 "payload",
1283 ArrowDataType::Int32,
1284 true,
1285 )])),
1286 true,
1287 )];
1288
1289 assert!(!schema_contains_same_fields(
1290 current_schema.fields(),
1291 &expected_fields
1292 ));
1293 }
1294
1295 #[test]
1296 fn test_count_snapshots_since_rewrite_in_metadata_ignores_other_branches() {
1297 let mut builder = TableMetadataBuilder::new(
1298 Schema::builder()
1299 .with_fields(vec![
1300 NestedField::new(1, "id", Type::Primitive(PrimitiveType::Long), false).into(),
1301 ])
1302 .build()
1303 .unwrap(),
1304 UnboundPartitionSpec::builder().build(),
1305 SortOrder::unsorted_order(),
1306 "s3://warehouse/db/table".to_owned(),
1307 FormatVersion::V2,
1308 HashMap::new(),
1309 )
1310 .unwrap();
1311
1312 for (snapshot_id, parent_snapshot_id, operation) in [
1313 (1, None, Operation::Append),
1314 (2, Some(1), Operation::Append),
1315 (3, Some(2), Operation::Replace),
1316 (4, Some(3), Operation::Append),
1317 (5, None, Operation::Overwrite),
1320 ] {
1321 builder = builder
1322 .add_snapshot(snapshot(snapshot_id, parent_snapshot_id, operation))
1323 .unwrap();
1324 }
1325
1326 let metadata = builder
1327 .set_ref(super::super::ICEBERG_COW_BRANCH, snapshot_ref(4))
1328 .unwrap()
1329 .set_ref(MAIN_BRANCH, snapshot_ref(5))
1330 .unwrap()
1331 .build()
1332 .unwrap()
1333 .metadata;
1334
1335 let count = IcebergSinkCommitter::count_snapshots_since_rewrite_in_metadata(
1336 &metadata,
1337 super::super::ICEBERG_COW_BRANCH,
1338 );
1339
1340 assert_eq!(count, 1);
1341 }
1342
1343 fn snapshot(
1344 snapshot_id: i64,
1345 parent_snapshot_id: Option<i64>,
1346 operation: Operation,
1347 ) -> Snapshot {
1348 Snapshot::builder()
1349 .with_snapshot_id(snapshot_id)
1350 .with_parent_snapshot_id(parent_snapshot_id)
1351 .with_sequence_number(snapshot_id)
1352 .with_timestamp_ms(snapshot_id)
1353 .with_manifest_list(format!("/snap-{snapshot_id}.avro"))
1354 .with_summary(Summary {
1355 operation,
1356 additional_properties: HashMap::new(),
1357 })
1358 .with_schema_id(0)
1359 .build()
1360 }
1361
1362 fn snapshot_ref(snapshot_id: i64) -> SnapshotReference {
1363 SnapshotReference::new(snapshot_id, SnapshotRetention::branch(None, None, None))
1364 }
1365}