Skip to main content

risingwave_connector/sink/iceberg/
commit.rs

1// Copyright 2026 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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        // RW Decimal has no precision/scale. Sink creation already accepts any
200        // table Decimal128 precision/scale, while expected schemas here are
201        // generated from RW columns using the same canonical mapping. Ignoring
202        // both values prevents false schema-state ambiguity and cannot mask an
203        // auto schema change, which only supports add/drop columns.
204        (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        // Commit data if present
356        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        // Deserialize commit metadata
489        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        // Last element is snapshot_id
497        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        // Remaining elements are write_results
507        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        // TODO: Files that have been written but not committed should be deleted.
575        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
586/// Methods Required to Achieve Exactly Once Semantics
587impl 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        // Skip if no data to commit
610        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        // guarantee that all write results has same schema_id and partition_spec_id
618        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        // Empty write results should be handled before calling this function.
654        assert!(
655            !write_results.is_empty() && !write_results.iter().all(|r| r.data_files.is_empty())
656        );
657
658        // Check snapshot limit before proceeding with commit
659        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        // Load the latest table to avoid concurrent modification with the best effort.
682        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        // # TODO:
711        // This retry behavior should be revert and do in iceberg-rust when it supports retry(Track in: https://github.com/apache/iceberg-rust/issues/964)
712        // because retry logic involved reapply the commit metadata.
713        // For now, we just retry the commit operation.
714        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    /// During pre-commit metadata, we record the `snapshot_id` corresponding to each batch of files.
823    /// Therefore, the logic for checking whether all files in this batch are present in Iceberg
824    /// has been changed to verifying if their corresponding `snapshot_id` exists in Iceberg.
825    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    /// Check if the specified columns already exist in the iceberg table's current schema.
839    /// This is used to determine if schema change has already been applied.
840    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 current schema equals original_schema, then schema change is NOT applied.
865        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 current schema equals the changed schema, then schema change is applied.
913        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    /// Commit schema changes (e.g., add columns) to the iceberg table.
928    /// This function uses Transaction API to atomically update the table schema
929    /// with optimistic locking to prevent concurrent conflicts.
930    async fn commit_schema_change_impl(&mut self, schema_change: PbSinkSchemaChange) -> Result<()> {
931        // Step 1: Build new fields to add
932        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                    // Convert RisingWave Field to Arrow Field using IcebergCreateTableArrowConvert
943                    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                    // Convert Arrow DataType to Iceberg Type
951                    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        // Step 2: Create Transaction with UpdateSchemaAction
977        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    /// Check the number of snapshots on the given branch lineage since the last rewrite operation.
1013    /// Returns the number of snapshots since the last rewrite.
1014    fn count_snapshots_since_rewrite_in_metadata(metadata: &TableMetadata, branch: &str) -> usize {
1015        // Start from the latest snapshot of the commit branch.
1016        let mut snapshot_id = metadata
1017            .snapshot_for_ref(branch)
1018            .map(|snapshot| snapshot.snapshot_id());
1019        let mut count = 0;
1020
1021        // Iterate through snapshots by parent lineage to find the last rewrite.
1022        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            // Check if this snapshot represents a rewrite operation.
1028            if snapshot.summary().operation == Operation::Replace {
1029                // Found a rewrite operation, stop counting.
1030                break;
1031            }
1032
1033            // Increment count for each snapshot that is not a rewrite.
1034            count += 1;
1035            snapshot_id = snapshot.parent_snapshot_id();
1036        }
1037
1038        count
1039    }
1040
1041    /// Returns the number of snapshots in the current commit branch since the last rewrite.
1042    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    /// Wait until snapshot count since last rewrite is below the limit
1048    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                // Wait for 30 seconds before checking again
1071                tokio::time::sleep(Duration::from_secs(30)).await;
1072
1073                // Refresh table after the wait so the next check sees latest snapshots.
1074                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            // Simulate the COW publish snapshot on main. It should not affect
1320            // the ingestion branch backlog count.
1321            (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}