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::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        // RW Decimal has no precision/scale. Sink creation already accepts any
199        // table Decimal128 precision/scale, while expected schemas here are
200        // generated from RW columns using the same canonical mapping. Ignoring
201        // both values prevents false schema-state ambiguity and cannot mask an
202        // auto schema change, which only supports add/drop columns.
203        (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        // Commit data if present
353        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        // Deserialize commit metadata
486        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        // Last element is snapshot_id
494        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        // Remaining elements are write_results
504        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        // TODO: Files that have been written but not committed should be deleted.
572        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
583/// Methods Required to Achieve Exactly Once Semantics
584impl 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        // Skip if no data to commit
607        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        // guarantee that all write results has same schema_id and partition_spec_id
615        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        // Empty write results should be handled before calling this function.
651        assert!(
652            !write_results.is_empty() && !write_results.iter().all(|r| r.data_files.is_empty())
653        );
654
655        // Check snapshot limit before proceeding with commit
656        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        // Load the latest table to avoid concurrent modification with the best effort.
679        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        // # TODO:
707        // 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)
708        // because retry logic involved reapply the commit metadata.
709        // For now, we just retry the commit operation.
710        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    /// During pre-commit metadata, we record the `snapshot_id` corresponding to each batch of files.
814    /// Therefore, the logic for checking whether all files in this batch are present in Iceberg
815    /// has been changed to verifying if their corresponding `snapshot_id` exists in Iceberg.
816    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    /// Check if the specified columns already exist in the iceberg table's current schema.
830    /// This is used to determine if schema change has already been applied.
831    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 current schema equals original_schema, then schema change is NOT applied.
856        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 current schema equals the changed schema, then schema change is applied.
904        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    /// Commit schema changes (e.g., add columns) to the iceberg table.
919    /// This function uses Transaction API to atomically update the table schema
920    /// with optimistic locking to prevent concurrent conflicts.
921    async fn commit_schema_change_impl(&mut self, schema_change: PbSinkSchemaChange) -> Result<()> {
922        use iceberg::spec::NestedField;
923
924        // Step 1: Get current table metadata
925        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        // Step 2: Build new fields to add
930        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                    // Convert RisingWave Field to Arrow Field using IcebergCreateTableArrowConvert
941                    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                    // Convert Arrow DataType to Iceberg Type
949                    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                    // Create NestedField with the next available field ID
958                    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        // Step 3: Create Transaction with UpdateSchemaAction
983        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    /// Check the number of snapshots on the given branch lineage since the last rewrite operation.
1016    /// Returns the number of snapshots since the last rewrite.
1017    fn count_snapshots_since_rewrite_in_metadata(metadata: &TableMetadata, branch: &str) -> usize {
1018        // Start from the latest snapshot of the commit branch.
1019        let mut snapshot_id = metadata
1020            .snapshot_for_ref(branch)
1021            .map(|snapshot| snapshot.snapshot_id());
1022        let mut count = 0;
1023
1024        // Iterate through snapshots by parent lineage to find the last rewrite.
1025        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            // Check if this snapshot represents a rewrite operation.
1031            if snapshot.summary().operation == Operation::Replace {
1032                // Found a rewrite operation, stop counting.
1033                break;
1034            }
1035
1036            // Increment count for each snapshot that is not a rewrite.
1037            count += 1;
1038            snapshot_id = snapshot.parent_snapshot_id();
1039        }
1040
1041        count
1042    }
1043
1044    /// Returns the number of snapshots in the current commit branch since the last rewrite.
1045    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    /// Wait until snapshot count since last rewrite is below the limit
1051    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                // Wait for 30 seconds before checking again
1074                tokio::time::sleep(Duration::from_secs(30)).await;
1075
1076                // Refresh table after the wait so the next check sees latest snapshots.
1077                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            // Simulate the COW publish snapshot on main. It should not affect
1323            // the ingestion branch backlog count.
1324            (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}