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