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 None,
688 )
689 .await
690 .map_err(SinkError::Iceberg)?;
691
692 let Some(schema) = self.table.metadata().schema_by_id(expect_schema_id) else {
693 return Err(SinkError::Iceberg(anyhow!(
694 "Can't find schema by id {}",
695 expect_schema_id
696 )));
697 };
698 let partition_type = resolve_partition_type(&self.table, expect_partition_spec_id, schema)?;
699
700 let data_files = write_results
701 .into_iter()
702 .flat_map(|r| {
703 r.data_files.into_iter().map(|f| {
704 f.try_into(expect_partition_spec_id, &partition_type, schema)
705 .map_err(|err| SinkError::Iceberg(anyhow!(err)))
706 })
707 })
708 .collect::<Result<Vec<DataFile>>>()?;
709
710 let catalog = self.catalog.clone();
715 let sink_id = self.sink_id;
716 let table_ident = self.table.identifier().clone();
717 let table_name = table_ident.to_string();
718 let retry_log_context = CommitRetryLogContext::new(
719 "sink_committer",
720 "commit",
721 table_name.clone(),
722 target_branch.clone(),
723 )
724 .with_sink_id(sink_id)
725 .with_epoch(epoch)
726 .with_snapshot_id(snapshot_id);
727
728 let table = commit_retry::run_with_retry(
729 catalog.clone(),
730 table_ident.clone(),
731 expect_schema_id,
732 expect_partition_spec_id,
733 None,
734 self.commit_retry_num as usize,
735 retry_log_context,
736 |table| {
737 let table_name = table_name.clone();
738 let target_branch = target_branch.clone();
739 let data_files = data_files.clone();
740 let catalog = catalog.clone();
741 async move {
742 let txn = Transaction::new(&table);
743 let append_action = txn
744 .fast_append()
745 .set_snapshot_id(snapshot_id)
746 .set_target_branch(target_branch.clone())
747 .set_snapshot_properties(HashMap::from([(
748 RISINGWAVE_ICEBERG_COMMIT_EPOCH.to_owned(),
749 epoch.to_string(),
750 )]))
751 .add_data_files(data_files);
752
753 let tx = append_action.apply(txn).map_err(|err| {
754 let err: IcebergError = err.into();
755 tracing::error!(
756 iceberg_component = "sink_committer",
757 iceberg_operation = "commit",
758 sink_id = %sink_id,
759 table = %table_name,
760 epoch,
761 snapshot_id,
762 branch = %target_branch,
763 data_file_count,
764 error = %err.as_report(),
765 "iceberg_sink_commit_fast_append_apply_failed",
766 );
767 CommitError::Commit(anyhow!(err).context("apply iceberg fast_append"))
768 })?;
769
770 let table = tx.commit(catalog.as_ref()).await.map_err(|err| {
771 let err: IcebergError = err.into();
772 tracing::error!(
773 iceberg_component = "sink_committer",
774 iceberg_operation = "commit",
775 sink_id = %sink_id,
776 table = %table_name,
777 epoch,
778 snapshot_id,
779 branch = %target_branch,
780 data_file_count,
781 error = %err.as_report(),
782 "iceberg_sink_commit_transaction_failed",
783 );
784 CommitError::Commit(anyhow!(err).context("commit iceberg transaction"))
785 })?;
786 Ok(table)
787 }
788 },
789 )
790 .await
791 .map_err(SinkError::Iceberg)?;
792 self.table = table;
793
794 let snapshot_num = self.table.metadata().snapshots().count();
795 let catalog_name = self.config.common.catalog_name();
796 let table_name = self.table.identifier().to_string();
797 let metrics_labels = [&self.param.sink_name, &catalog_name, &table_name];
798 GLOBAL_SINK_METRICS
799 .iceberg_snapshot_num
800 .with_guarded_label_values(&metrics_labels)
801 .set(snapshot_num as i64);
802
803 tracing::debug!(
804 iceberg_component = "sink_committer",
805 iceberg_operation = "commit",
806 sink_id = %self.sink_id,
807 sink_name = %self.param.sink_name,
808 table = %self.table.identifier(),
809 branch = %target_branch,
810 epoch,
811 snapshot_id,
812 snapshot_num,
813 data_file_count,
814 "iceberg_sink_commit_succeeded",
815 );
816
817 self.notify_iceberg_compaction_scheduler(false);
818
819 Ok(())
820 }
821
822 async fn is_snapshot_id_in_iceberg(&self, snapshot_id: i64) -> Result<bool> {
826 let table = self
827 .catalog
828 .load_table(self.table.identifier())
829 .await
830 .map_err(|err| SinkError::Iceberg(anyhow!(err).context("reload iceberg table")))?;
831 if table.metadata().snapshot_by_id(snapshot_id).is_some() {
832 Ok(true)
833 } else {
834 Ok(false)
835 }
836 }
837
838 fn check_schema_change_applied(&self, schema_change: &PbSinkSchemaChange) -> Result<bool> {
841 let current_schema = self.table.metadata().current_schema();
842 let current_arrow_schema = schema_to_arrow_schema(current_schema.as_ref())
843 .context("Failed to convert schema")
844 .map_err(SinkError::Iceberg)?;
845
846 let iceberg_arrow_convert = IcebergArrowConvert;
847
848 let schema_matches = |expected: &[ArrowField]| {
849 schema_contains_same_fields(current_arrow_schema.fields(), expected)
850 };
851
852 let original_arrow_fields: Vec<ArrowField> = schema_change
853 .original_schema
854 .iter()
855 .map(|pb_field| {
856 let field = Field::from(pb_field);
857 iceberg_arrow_convert
858 .to_arrow_field(&field.name, &field.data_type)
859 .context("Failed to convert field to arrow")
860 .map_err(SinkError::Iceberg)
861 })
862 .collect::<Result<_>>()?;
863
864 if schema_matches(&original_arrow_fields) {
866 tracing::debug!(
867 "Current iceberg schema matches original_schema ({} columns); schema change not applied",
868 original_arrow_fields.len()
869 );
870 return Ok(false);
871 }
872
873 let expected_after_change = match schema_change.op.as_ref() {
874 Some(risingwave_pb::stream_plan::sink_schema_change::Op::AddColumns(
875 add_columns_op,
876 )) => {
877 let add_arrow_fields: Vec<ArrowField> = add_columns_op
878 .fields
879 .iter()
880 .map(|pb_field| {
881 let field = Field::from(pb_field);
882 iceberg_arrow_convert
883 .to_arrow_field(&field.name, &field.data_type)
884 .context("Failed to convert field to arrow")
885 .map_err(SinkError::Iceberg)
886 })
887 .collect::<Result<_>>()?;
888
889 let mut expected_after_change = original_arrow_fields;
890 expected_after_change.extend(add_arrow_fields);
891 expected_after_change
892 }
893 Some(risingwave_pb::stream_plan::sink_schema_change::Op::DropColumns(
894 drop_columns_op,
895 )) => original_arrow_fields
896 .into_iter()
897 .filter(|field| {
898 !drop_columns_op
899 .column_names
900 .iter()
901 .any(|name| name == field.name())
902 })
903 .collect_vec(),
904 _ => {
905 return Err(SinkError::Iceberg(anyhow!(
906 "Unsupported sink schema change op in iceberg sink: {:?}",
907 schema_change.op
908 )));
909 }
910 };
911
912 if schema_matches(&expected_after_change) {
914 tracing::debug!(
915 "Current iceberg schema matches changed schema ({} columns); schema change already applied",
916 expected_after_change.len()
917 );
918 return Ok(true);
919 }
920
921 Err(SinkError::Iceberg(anyhow!(
922 "Current iceberg schema does not match either original_schema ({} cols) or changed schema; cannot determine whether schema change is applied",
923 schema_change.original_schema.len()
924 )))
925 }
926
927 async fn commit_schema_change_impl(&mut self, schema_change: PbSinkSchemaChange) -> Result<()> {
931 let iceberg_create_table_arrow_convert = IcebergCreateTableArrowConvert::default();
933 let mut new_fields = Vec::new();
934
935 let mut drop_column_names = Vec::new();
936 match schema_change.op.as_ref() {
937 Some(risingwave_pb::stream_plan::sink_schema_change::Op::AddColumns(
938 add_columns_op,
939 )) => {
940 let add_columns = add_columns_op.fields.iter().map(Field::from).collect_vec();
941 for field in &add_columns {
942 let arrow_field = iceberg_create_table_arrow_convert
944 .to_arrow_field(&field.name, &field.data_type)
945 .with_context(|| {
946 format!("Failed to convert field '{}' to arrow", field.name)
947 })
948 .map_err(SinkError::Iceberg)?;
949
950 let iceberg_type = iceberg::arrow::arrow_type_to_type(arrow_field.data_type())
952 .map_err(|err| {
953 SinkError::Iceberg(
954 anyhow!(err)
955 .context("Failed to convert Arrow type to Iceberg type"),
956 )
957 })?;
958
959 new_fields.push(AddColumn::optional(&field.name, iceberg_type));
960 tracing::info!("Prepared field '{}' for schema change", field.name);
961 }
962 }
963 Some(risingwave_pb::stream_plan::sink_schema_change::Op::DropColumns(
964 drop_columns_op,
965 )) => {
966 drop_column_names = drop_columns_op.column_names.clone();
967 }
968 _ => {
969 return Err(SinkError::Iceberg(anyhow!(
970 "Unsupported sink schema change op in iceberg sink: {:?}",
971 schema_change.op
972 )));
973 }
974 }
975
976 tracing::info!(
978 "Committing schema change to catalog for table {}",
979 self.table.identifier()
980 );
981
982 let txn = Transaction::new(&self.table);
983 let action_fields_added = new_fields.len();
984 let mut action = txn.update_schema();
985 for field in new_fields {
986 action = action.add_column(field);
987 }
988 for column_name in &drop_column_names {
989 action = action.delete_column(column_name);
990 }
991
992 let updated_table = action
993 .apply(txn)
994 .context("Failed to apply schema update action")
995 .map_err(SinkError::Iceberg)?
996 .commit(self.catalog.as_ref())
997 .await
998 .context("Failed to commit table schema change")
999 .map_err(SinkError::Iceberg)?;
1000
1001 self.table = updated_table;
1002
1003 tracing::info!(
1004 "Successfully committed schema change, added {} columns and dropped {} columns from iceberg table",
1005 action_fields_added,
1006 drop_column_names.len()
1007 );
1008
1009 Ok(())
1010 }
1011
1012 fn count_snapshots_since_rewrite_in_metadata(metadata: &TableMetadata, branch: &str) -> usize {
1015 let mut snapshot_id = metadata
1017 .snapshot_for_ref(branch)
1018 .map(|snapshot| snapshot.snapshot_id());
1019 let mut count = 0;
1020
1021 while let Some(current_snapshot_id) = snapshot_id {
1023 let Some(snapshot) = metadata.snapshot_by_id(current_snapshot_id) else {
1024 break;
1025 };
1026
1027 if snapshot.summary().operation == Operation::Replace {
1029 break;
1031 }
1032
1033 count += 1;
1035 snapshot_id = snapshot.parent_snapshot_id();
1036 }
1037
1038 count
1039 }
1040
1041 fn count_snapshots_since_rewrite(&self) -> usize {
1043 let branch = commit_branch(self.config.r#type.as_str(), self.config.write_mode);
1044 Self::count_snapshots_since_rewrite_in_metadata(self.table.metadata(), branch.as_str())
1045 }
1046
1047 async fn wait_for_snapshot_limit(&mut self) -> Result<()> {
1049 if let Some(max_snapshots) = self.config.max_snapshots_num_before_compaction {
1050 loop {
1051 let current_count = self.count_snapshots_since_rewrite();
1052
1053 if current_count < max_snapshots {
1054 tracing::info!(
1055 "Snapshot count check passed: {} < {}",
1056 current_count,
1057 max_snapshots
1058 );
1059 break;
1060 }
1061
1062 tracing::info!(
1063 "Snapshot count {} exceeds limit {}, waiting...",
1064 current_count,
1065 max_snapshots
1066 );
1067
1068 self.notify_iceberg_compaction_scheduler(true);
1069
1070 tokio::time::sleep(Duration::from_secs(30)).await;
1072
1073 let table_ident = self.table.identifier().clone();
1075 self.table = self.catalog.load_table(&table_ident).await.map_err(|err| {
1076 SinkError::Iceberg(anyhow!(err).context("reload iceberg table"))
1077 })?;
1078 }
1079 }
1080 Ok(())
1081 }
1082}
1083
1084fn serialize_metadata(metadata: Vec<Vec<u8>>) -> Vec<u8> {
1085 serde_json::to_vec(&metadata).unwrap()
1086}
1087
1088fn deserialize_metadata(bytes: Vec<u8>) -> Vec<Vec<u8>> {
1089 serde_json::from_slice(&bytes).unwrap()
1090}
1091
1092#[cfg(test)]
1093mod tests {
1094 use std::collections::HashMap;
1095
1096 use iceberg::spec::{
1097 FormatVersion, MAIN_BRANCH, NestedField, PrimitiveType, Schema, Snapshot,
1098 SnapshotReference, SnapshotRetention, SortOrder, Summary, TableMetadataBuilder, Type,
1099 UnboundPartitionSpec,
1100 };
1101 use risingwave_common::array::arrow::arrow_schema_iceberg::{
1102 DataType as ArrowDataType, Field as ArrowField, FieldRef as ArrowFieldRef,
1103 Fields as ArrowFields, Schema as ArrowSchema,
1104 };
1105
1106 use super::*;
1107
1108 #[test]
1109 fn test_schema_contains_same_fields_allows_binary_large_binary() {
1110 let current_schema = ArrowSchema::new(vec![
1111 ArrowField::new("k", ArrowDataType::Int32, true),
1112 ArrowField::new("v", ArrowDataType::LargeBinary, true),
1113 ]);
1114 let expected_fields = vec![
1115 ArrowField::new("k", ArrowDataType::Int32, true),
1116 ArrowField::new("v", ArrowDataType::Binary, true),
1117 ];
1118
1119 assert!(schema_contains_same_fields(
1120 current_schema.fields(),
1121 &expected_fields
1122 ));
1123 }
1124
1125 #[test]
1126 fn test_schema_contains_same_fields_allows_nested_binary_large_binary() {
1127 let current_schema = ArrowSchema::new(vec![
1128 ArrowField::new(
1129 "s",
1130 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1131 "payload",
1132 ArrowDataType::LargeBinary,
1133 true,
1134 )])),
1135 true,
1136 ),
1137 ArrowField::new(
1138 "l",
1139 ArrowDataType::List(ArrowFieldRef::new(ArrowField::new_list_field(
1140 ArrowDataType::LargeBinary,
1141 true,
1142 ))),
1143 true,
1144 ),
1145 ArrowField::new_map(
1146 "m",
1147 "entries",
1148 ArrowFieldRef::new(ArrowField::new("key", ArrowDataType::Utf8, false)),
1149 ArrowFieldRef::new(ArrowField::new("value", ArrowDataType::LargeBinary, true)),
1150 false,
1151 true,
1152 ),
1153 ]);
1154 let expected_fields = vec![
1155 ArrowField::new(
1156 "s",
1157 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1158 "payload",
1159 ArrowDataType::Binary,
1160 true,
1161 )])),
1162 true,
1163 ),
1164 ArrowField::new(
1165 "l",
1166 ArrowDataType::List(ArrowFieldRef::new(ArrowField::new_list_field(
1167 ArrowDataType::Binary,
1168 true,
1169 ))),
1170 true,
1171 ),
1172 ArrowField::new_map(
1173 "m",
1174 "entries",
1175 ArrowFieldRef::new(ArrowField::new("key", ArrowDataType::Utf8, false)),
1176 ArrowFieldRef::new(ArrowField::new("value", ArrowDataType::Binary, true)),
1177 false,
1178 true,
1179 ),
1180 ];
1181
1182 assert!(schema_contains_same_fields(
1183 current_schema.fields(),
1184 &expected_fields
1185 ));
1186 }
1187
1188 #[test]
1189 fn test_schema_contains_same_fields_allows_decimal_precision_delta() {
1190 let current_schema = ArrowSchema::new(vec![ArrowField::new(
1191 "d",
1192 ArrowDataType::Decimal128(28, 10),
1193 true,
1194 )]);
1195 let expected_fields = vec![ArrowField::new(
1196 "d",
1197 ArrowDataType::Decimal128(38, 10),
1198 true,
1199 )];
1200
1201 assert!(schema_contains_same_fields(
1202 current_schema.fields(),
1203 &expected_fields
1204 ));
1205 }
1206
1207 #[test]
1208 fn test_schema_contains_same_fields_allows_decimal_precision_and_scale_delta() {
1209 let current_schema = ArrowSchema::new(vec![ArrowField::new(
1210 "d",
1211 ArrowDataType::Decimal128(38, 2),
1212 true,
1213 )]);
1214 let expected_fields = vec![ArrowField::new(
1215 "d",
1216 ArrowDataType::Decimal128(38, 10),
1217 true,
1218 )];
1219
1220 assert!(schema_contains_same_fields(
1221 current_schema.fields(),
1222 &expected_fields
1223 ));
1224 }
1225
1226 #[test]
1227 fn test_schema_contains_same_fields_rejects_length_mismatch() {
1228 let current_schema =
1229 ArrowSchema::new(vec![ArrowField::new("k", ArrowDataType::Int32, true)]);
1230 let expected_fields = vec![
1231 ArrowField::new("k", ArrowDataType::Int32, true),
1232 ArrowField::new("v", ArrowDataType::Utf8, true),
1233 ];
1234
1235 assert!(!schema_contains_same_fields(
1236 current_schema.fields(),
1237 &expected_fields
1238 ));
1239 }
1240
1241 #[test]
1242 fn test_schema_contains_same_fields_rejects_type_mismatch() {
1243 let current_schema =
1244 ArrowSchema::new(vec![ArrowField::new("v", ArrowDataType::Utf8, true)]);
1245 let expected_fields = vec![ArrowField::new("v", ArrowDataType::Int32, true)];
1246
1247 assert!(!schema_contains_same_fields(
1248 current_schema.fields(),
1249 &expected_fields
1250 ));
1251 }
1252
1253 #[test]
1254 fn test_schema_contains_same_fields_rejects_reused_duplicate_match() {
1255 let current_schema = ArrowSchema::new(vec![
1256 ArrowField::new("v", ArrowDataType::Int32, true),
1257 ArrowField::new("v", ArrowDataType::Utf8, true),
1258 ]);
1259 let expected_fields = vec![
1260 ArrowField::new("v", ArrowDataType::Int32, true),
1261 ArrowField::new("v", ArrowDataType::Int32, true),
1262 ];
1263
1264 assert!(!schema_contains_same_fields(
1265 current_schema.fields(),
1266 &expected_fields
1267 ));
1268 }
1269
1270 #[test]
1271 fn test_schema_contains_same_fields_rejects_nested_type_mismatch() {
1272 let current_schema = ArrowSchema::new(vec![ArrowField::new(
1273 "s",
1274 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1275 "payload",
1276 ArrowDataType::Utf8,
1277 true,
1278 )])),
1279 true,
1280 )]);
1281 let expected_fields = vec![ArrowField::new(
1282 "s",
1283 ArrowDataType::Struct(ArrowFields::from(vec![ArrowField::new(
1284 "payload",
1285 ArrowDataType::Int32,
1286 true,
1287 )])),
1288 true,
1289 )];
1290
1291 assert!(!schema_contains_same_fields(
1292 current_schema.fields(),
1293 &expected_fields
1294 ));
1295 }
1296
1297 #[test]
1298 fn test_count_snapshots_since_rewrite_in_metadata_ignores_other_branches() {
1299 let mut builder = TableMetadataBuilder::new(
1300 Schema::builder()
1301 .with_fields(vec![
1302 NestedField::new(1, "id", Type::Primitive(PrimitiveType::Long), false).into(),
1303 ])
1304 .build()
1305 .unwrap(),
1306 UnboundPartitionSpec::builder().build(),
1307 SortOrder::unsorted_order(),
1308 "s3://warehouse/db/table".to_owned(),
1309 FormatVersion::V2,
1310 HashMap::new(),
1311 )
1312 .unwrap();
1313
1314 for (snapshot_id, parent_snapshot_id, operation) in [
1315 (1, None, Operation::Append),
1316 (2, Some(1), Operation::Append),
1317 (3, Some(2), Operation::Replace),
1318 (4, Some(3), Operation::Append),
1319 (5, None, Operation::Overwrite),
1322 ] {
1323 builder = builder
1324 .add_snapshot(snapshot(snapshot_id, parent_snapshot_id, operation))
1325 .unwrap();
1326 }
1327
1328 let metadata = builder
1329 .set_ref(super::super::ICEBERG_COW_BRANCH, snapshot_ref(4))
1330 .unwrap()
1331 .set_ref(MAIN_BRANCH, snapshot_ref(5))
1332 .unwrap()
1333 .build()
1334 .unwrap()
1335 .metadata;
1336
1337 let count = IcebergSinkCommitter::count_snapshots_since_rewrite_in_metadata(
1338 &metadata,
1339 super::super::ICEBERG_COW_BRANCH,
1340 );
1341
1342 assert_eq!(count, 1);
1343 }
1344
1345 fn snapshot(
1346 snapshot_id: i64,
1347 parent_snapshot_id: Option<i64>,
1348 operation: Operation,
1349 ) -> Snapshot {
1350 Snapshot::builder()
1351 .with_snapshot_id(snapshot_id)
1352 .with_parent_snapshot_id(parent_snapshot_id)
1353 .with_sequence_number(snapshot_id)
1354 .with_timestamp_ms(snapshot_id)
1355 .with_manifest_list(format!("/snap-{snapshot_id}.avro"))
1356 .with_summary(Summary {
1357 operation,
1358 additional_properties: HashMap::new(),
1359 })
1360 .with_schema_id(0)
1361 .build()
1362 }
1363
1364 fn snapshot_ref(snapshot_id: i64) -> SnapshotReference {
1365 SnapshotReference::new(snapshot_id, SnapshotRetention::branch(None, None, None))
1366 }
1367}