Skip to main content

mz_storage/sink/
iceberg.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Iceberg sink implementation.
11//!
12//! This code renders a [`IcebergSinkConnection`] into a dataflow that writes
13//! data to an Iceberg table. `SinkRender::render_sink` hands the sink a stream
14//! of arrangement batches keyed on the sink key (the upstream arrangement's
15//! trace reader is already dropped, so the spine is free to compact as
16//! batches flow). A small `walk_sink_arrangement` operator consumes that
17//! stream and emits one `DiffPair` per `(key, timestamp)` update into the
18//! pipeline below.
19//!
20//! ```text
21//!        ┏━━━━━━━━━━━━━━┓
22//!        ┃   persist    ┃
23//!        ┃    source    ┃
24//!        ┗━━━━━━┯━━━━━━━┛
25//!               │ stream of arrangement batches (trace reader dropped)
26//!               │
27//!        ┏━━━━━━v━━━━━━━┓
28//!        ┃    walk      ┃
29//!        ┃ arrangement  ┃ yields individual DiffPairs per (key, timestamp)
30//!        ┗━━━━━━┯━━━━━━━┛
31//!               │ (Option<Row>, DiffPair<Row>) rows
32//!               │
33//!        ┏━━━━━━v━━━━━━━┓
34//!        ┃     mint     ┃ (single worker)
35//!        ┃    batch     ┃ loads/creates the Iceberg table,
36//!        ┃ descriptions ┃ determines resume upper
37//!        ┗━━━┯━━━━━┯━━━━┛
38//!            │     │ batch descriptions (broadcast)
39//!       rows │     ├─────────────────────────┐
40//!            │     │                         │
41//!        ┏━━━v━━━━━v━━━━┓    ╭─────────────╮ │
42//!        ┃    write     ┃───>│ S3 / object │ │
43//!        ┃  data files  ┃    │   storage   │ │
44//!        ┗━━━━━━┯━━━━━━━┛    ╰─────────────╯ │
45//!               │ file metadata              │
46//!               │                            │
47//!        ┏━━━━━━v━━━━━━━━━━━━━━━━━━━━━━━━━━━━v┓
48//!        ┃           commit to                ┃ (single worker)
49//!        ┃             iceberg                ┃
50//!        ┗━━━━━━━━━━━━━┯━━━━━━━━━━━━━━━━━━━━━━┛
51//!                      │
52//!              ╭───────v───────╮
53//!              │ Iceberg table │
54//!              │  (snapshots)  │
55//!              ╰───────────────╯
56//! ```
57//! # Minting batch descriptions
58//! The "mint batch descriptions" operator is responsible for generating
59//! time-based batch boundaries that group writes into Iceberg snapshots.
60//! It maintains a sliding window of future batch descriptions so that
61//! writers can start processing data even while earlier batches are still being written.
62//! Knowing the batch boundaries ahead of time is important because we need to
63//! be able to make the claim that all data files written for a given batch
64//! include all data up to the upper `t` but not beyond it.
65//! This could be trivially achieved by waiting for all data to arrive up to a certain
66//! frontier, but that would prevent us from streaming writes out to object storage
67//! until the entire batch is complete, which would increase latency and reduce throughput.
68//!
69//! # Writing data files
70//! The "write data files" operator receives rows along with batch descriptions.
71//! It matches rows to batches by timestamp; if a batch description hasn't arrived yet,
72//! rows are stashed until it does. This allows batches to be minted ahead of data arrival.
73//! The operator uses an Iceberg `DeltaWriter` to write Parquet data files
74//! (and position delete files if necessary) to object storage.
75//! It outputs metadata about the written files along with their batch descriptions
76//! for the commit operator to consume.
77//!
78//! # Committing to Iceberg
79//! The "commit to iceberg" operator receives metadata about written data files
80//! along with their batch descriptions. It groups files by batch and creates
81//! Iceberg snapshots that include all files for each batch. It updates the Iceberg
82//! table's metadata to reflect the new snapshots, including updating the
83//! `mz-frontier` property to track progress.
84
85use std::cmp::Ordering;
86use std::collections::VecDeque;
87use std::convert::Infallible;
88use std::future::Future;
89use std::time::Instant;
90use std::{cell::RefCell, rc::Rc, sync::Arc};
91
92use anyhow::{Context, anyhow};
93use arrow::array::{ArrayRef, Int32Array, Int64Array, RecordBatch};
94use arrow::datatypes::{DataType, Field, Schema as ArrowSchema};
95use differential_dataflow::trace::BatchReader;
96use differential_dataflow::trace::implementations::ord_neu::OrdValBatch;
97use differential_dataflow::trace::implementations::{BatchContainer, Layout};
98use differential_dataflow::{Hashable, VecCollection};
99use futures::StreamExt;
100use iceberg::ErrorKind;
101use iceberg::arrow::{arrow_schema_to_schema, schema_to_arrow_schema};
102use iceberg::spec::{
103    DataFile, FormatVersion, NestedField, PrimitiveType, Snapshot, StructType, Type,
104    read_data_files_from_avro, write_data_files_to_avro,
105};
106use iceberg::spec::{Schema, SchemaRef};
107use iceberg::table::Table;
108use iceberg::transaction::{RowDeltaAction, TransactionAction};
109use iceberg::writer::base_writer::data_file_writer::DataFileWriterBuilder;
110use iceberg::writer::base_writer::equality_delete_writer::{
111    EqualityDeleteFileWriterBuilder, EqualityDeleteWriterConfig,
112};
113use iceberg::writer::base_writer::position_delete_writer::{
114    PositionDeleteFileWriterBuilder, PositionDeleteWriterConfig,
115};
116use iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder;
117use iceberg::writer::file_writer::ParquetWriterBuilder;
118use iceberg::writer::file_writer::location_generator::{
119    DefaultFileNameGenerator, DefaultLocationGenerator,
120};
121use iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder;
122use iceberg::writer::{IcebergWriter, IcebergWriterBuilder};
123use iceberg::{Catalog, NamespaceIdent, TableCommit, TableCreation, TableIdent};
124use itertools::Itertools;
125use mz_arrow_util::builder::{ARROW_EXTENSION_NAME_KEY, ArrowBuilder};
126use mz_interchange::avro::DiffPair;
127use mz_interchange::envelopes::for_each_diff_pair_async;
128use mz_ore::cast::CastFrom;
129use mz_ore::error::ErrorExt;
130use mz_ore::future::InTask;
131use mz_ore::result::ResultExt;
132use mz_ore::retry::{Retry, RetryResult};
133use mz_persist_client::Diagnostics;
134use mz_persist_client::write::WriteHandle;
135use mz_persist_types::codec_impls::UnitSchema;
136use mz_repr::{Diff, GlobalId, Row, Timestamp};
137use mz_row_spine::ArcBatch;
138use mz_storage_types::StorageDiff;
139use mz_storage_types::configuration::StorageConfiguration;
140use mz_storage_types::controller::CollectionMetadata;
141use mz_storage_types::errors::DataflowError;
142use mz_storage_types::sinks::{
143    IcebergSinkConnection, SinkEnvelope, StorageSinkDesc, iceberg_type_overrides,
144};
145use mz_storage_types::sources::SourceData;
146use mz_timely_util::antichain::AntichainExt;
147use mz_timely_util::builder_async::{Event, OperatorBuilder, PressOnDropButton};
148use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
149use parquet::file::properties::WriterProperties;
150use serde::{Deserialize, Serialize};
151use timely::PartialOrder;
152use timely::container::CapacityContainerBuilder;
153use timely::dataflow::StreamVec;
154use timely::dataflow::channels::pact::{Exchange, Pipeline};
155use timely::dataflow::operators::vec::{Broadcast, Map, ToStream};
156use timely::dataflow::operators::{CapabilitySet, Concatenate};
157use timely::progress::{Antichain, Timestamp as _};
158use tracing::{debug, info};
159
160use crate::healthcheck::{HealthStatusMessage, HealthStatusUpdate, StatusNamespace};
161use crate::metrics::sink::iceberg::IcebergSinkMetrics;
162use crate::render::sinks::{PkViolationWarner, SinkBatchStream, SinkRender};
163use crate::statistics::SinkStatistics;
164use crate::storage_state::StorageState;
165
166/// Set the default capacity for the array builders inside the ArrowBuilder. This is the
167/// number of items each builder can hold before it needs to allocate more memory.
168const DEFAULT_ARRAY_BUILDER_ITEM_CAPACITY: usize = 1024;
169/// Set the default buffer capacity for the string and binary array builders inside the
170/// ArrowBuilder. This is the number of bytes each builder can hold before it needs to allocate
171/// more memory.
172const DEFAULT_ARRAY_BUILDER_DATA_CAPACITY: usize = 1024;
173
174/// The prefix for Parquet files written by this sink.
175const PARQUET_FILE_PREFIX: &str = "mz_data";
176/// The number of batch descriptions to mint ahead of the observed frontier. This determines how
177/// many batches we have in-flight at any given time.
178const INITIAL_DESCRIPTIONS_TO_MINT: u64 = 3;
179
180/// Shared state produced by the async setup in [`write_data_files`] that both
181/// envelope handlers need to construct Parquet writers.
182struct WriterContext {
183    /// Arrow schema for data columns, with Materialize extension metadata merged in.
184    arrow_schema: Arc<ArrowSchema>,
185    /// Iceberg table schema, used to configure Parquet writers.
186    current_schema: Arc<Schema>,
187    /// File I/O for writing Parquet files to object storage.
188    file_io: iceberg::io::FileIO,
189    /// Generates file paths under the table's data directory.
190    location_generator: DefaultLocationGenerator,
191    /// Generates unique file names with a per-worker UUID suffix.
192    file_name_generator: DefaultFileNameGenerator,
193    writer_properties: WriterProperties,
194}
195
196/// Envelope-specific logic for writing Iceberg data files.
197trait EnvelopeHandler: Send {
198    /// Construct from the shared writer context after async setup completes.
199    fn new(
200        ctx: WriterContext,
201        connection: &IcebergSinkConnection,
202        materialize_arrow_schema: &Arc<ArrowSchema>,
203    ) -> anyhow::Result<Self>
204    where
205        Self: Sized;
206
207    /// Create an [`IcebergWriter`] for a new batch.
208    ///
209    /// `is_snapshot` is true for the initial "snapshot" batch (lower == as_of), which
210    /// contains all pre-existing data and can be very large. Implementations may use
211    /// this to disable memory-intensive optimisations like seen-rows deduplication.
212    async fn create_writer(&self, is_snapshot: bool) -> anyhow::Result<Box<dyn IcebergWriter>>;
213
214    fn row_to_batch(&self, diff_pair: DiffPair<Row>, ts: Timestamp) -> anyhow::Result<RecordBatch>;
215}
216
217struct UpsertEnvelopeHandler {
218    ctx: WriterContext,
219    /// Iceberg field IDs of the key columns, used for equality delete files.
220    equality_ids: Vec<i32>,
221    /// Iceberg schema for position delete files.
222    pos_schema: Arc<Schema>,
223    /// Iceberg schema for equality delete files (projected to key columns only).
224    eq_schema: Arc<Schema>,
225    /// Configuration for the equality delete writer (projected schema + column IDs).
226    eq_config: EqualityDeleteWriterConfig,
227    /// Arrow schema with an appended `__op` column that the
228    /// [`DeltaWriter`](iceberg::writer::combined_writer::delta_writer::DeltaWriter)
229    /// uses to distinguish inserts (+1) from deletes (-1).
230    schema_with_op: Arc<ArrowSchema>,
231}
232
233impl EnvelopeHandler for UpsertEnvelopeHandler {
234    fn new(
235        ctx: WriterContext,
236        connection: &IcebergSinkConnection,
237        materialize_arrow_schema: &Arc<ArrowSchema>,
238    ) -> anyhow::Result<Self> {
239        let Some((_, equality_indices)) = &connection.key_desc_and_indices else {
240            return Err(anyhow::anyhow!(
241                "Iceberg sink requires key columns for equality deletes"
242            ));
243        };
244
245        let equality_ids = equality_ids_for_indices(
246            ctx.current_schema.as_ref(),
247            materialize_arrow_schema.as_ref(),
248            equality_indices,
249        )?;
250
251        let pos_arrow_schema = PositionDeleteWriterConfig::arrow_schema();
252        let pos_schema = Arc::new(
253            arrow_schema_to_schema(&pos_arrow_schema)
254                .context("Failed to convert position delete Arrow schema to Iceberg schema")?,
255        );
256
257        let eq_config =
258            EqualityDeleteWriterConfig::new(equality_ids.clone(), Arc::clone(&ctx.current_schema))
259                .context("Failed to create EqualityDeleteWriterConfig")?;
260        let eq_schema = Arc::new(
261            arrow_schema_to_schema(eq_config.projected_arrow_schema_ref())
262                .context("Failed to convert equality delete Arrow schema to Iceberg schema")?,
263        );
264
265        let schema_with_op = Arc::new(build_schema_with_op_column(&ctx.arrow_schema));
266
267        Ok(Self {
268            ctx,
269            equality_ids,
270            pos_schema,
271            eq_schema,
272            eq_config,
273            schema_with_op,
274        })
275    }
276
277    async fn create_writer(&self, is_snapshot: bool) -> anyhow::Result<Box<dyn IcebergWriter>> {
278        let data_parquet_writer = ParquetWriterBuilder::new(
279            self.ctx.writer_properties.clone(),
280            Arc::clone(&self.ctx.current_schema),
281        )
282        .with_arrow_schema(Arc::clone(&self.ctx.arrow_schema))
283        .context("Arrow schema validation failed")?;
284        let data_rolling_writer = RollingFileWriterBuilder::new_with_default_file_size(
285            data_parquet_writer,
286            Arc::clone(&self.ctx.current_schema),
287            self.ctx.file_io.clone(),
288            self.ctx.location_generator.clone(),
289            self.ctx.file_name_generator.clone(),
290        );
291        let data_writer_builder = DataFileWriterBuilder::new(data_rolling_writer);
292
293        let pos_config = PositionDeleteWriterConfig::new(None, 0, None);
294        let pos_parquet_writer = ParquetWriterBuilder::new(
295            self.ctx.writer_properties.clone(),
296            Arc::clone(&self.pos_schema),
297        );
298        let pos_rolling_writer = RollingFileWriterBuilder::new_with_default_file_size(
299            pos_parquet_writer,
300            Arc::clone(&self.ctx.current_schema),
301            self.ctx.file_io.clone(),
302            self.ctx.location_generator.clone(),
303            self.ctx.file_name_generator.clone(),
304        );
305        let pos_delete_writer_builder =
306            PositionDeleteFileWriterBuilder::new(pos_rolling_writer, pos_config);
307
308        let eq_parquet_writer = ParquetWriterBuilder::new(
309            self.ctx.writer_properties.clone(),
310            Arc::clone(&self.eq_schema),
311        );
312        let eq_rolling_writer = RollingFileWriterBuilder::new_with_default_file_size(
313            eq_parquet_writer,
314            Arc::clone(&self.ctx.current_schema),
315            self.ctx.file_io.clone(),
316            self.ctx.location_generator.clone(),
317            self.ctx.file_name_generator.clone(),
318        );
319        let eq_delete_writer_builder =
320            EqualityDeleteFileWriterBuilder::new(eq_rolling_writer, self.eq_config.clone());
321
322        let mut builder = DeltaWriterBuilder::new(
323            data_writer_builder,
324            pos_delete_writer_builder,
325            eq_delete_writer_builder,
326            self.equality_ids.clone(),
327        );
328
329        builder = if is_snapshot {
330            // Snapshot batches only produce inserts, so disable seen_rows tracking to save memory.
331            builder.with_max_seen_rows(0)
332        } else {
333            // For incremental batches, keep all "seen" rows. Do not evict any rows.
334            // The DeltaWriter issues an equality delete if we update (or delete) a row outside the "seen" cache.
335            // But equality deletes only apply to prior snapshots (lower sequence number).
336            //
337            // i.e. The DeltaWriter assumes that rows outside the "seen" cache come from prior snapshots.
338            //
339            // If we insert a row a=foo during this snapshot and then evict it from the "seen" cache,
340            // a subsequent update a=bar (also during this snapshot) will lead to:
341            //   1. Equality delete for a=foo (does nothing because a=foo is from this snapshot, not a prior snapshot)
342            //   2. Insert a=bar
343            // Because the deletion does nothing, we have a=foo and a=bar in the same snapshot.
344            builder.with_max_seen_rows(usize::MAX)
345        };
346
347        Ok(Box::new(
348            builder
349                .build(None)
350                .await
351                .context("Failed to create DeltaWriter")?,
352        ))
353    }
354
355    /// The `__op` column indicates whether each row is an insert (+1) or delete (-1),
356    /// which the DeltaWriter uses to generate the appropriate Iceberg data/delete files.
357    fn row_to_batch(
358        &self,
359        diff_pair: DiffPair<Row>,
360        _ts: Timestamp,
361    ) -> anyhow::Result<RecordBatch> {
362        let mut builder = ArrowBuilder::new_with_schema(
363            Arc::clone(&self.ctx.arrow_schema),
364            DEFAULT_ARRAY_BUILDER_ITEM_CAPACITY,
365            DEFAULT_ARRAY_BUILDER_DATA_CAPACITY,
366        )
367        .context("Failed to create builder")?;
368
369        let mut op_values = Vec::new();
370
371        if let Some(before) = diff_pair.before {
372            builder
373                .add_row(&before)
374                .context("Failed to add delete row to builder")?;
375            op_values.push(-1i32);
376        }
377        if let Some(after) = diff_pair.after {
378            builder
379                .add_row(&after)
380                .context("Failed to add insert row to builder")?;
381            op_values.push(1i32);
382        }
383
384        let batch = builder
385            .to_record_batch()
386            .context("Failed to create record batch")?;
387
388        let mut columns: Vec<ArrayRef> = batch.columns().to_vec();
389        columns.push(Arc::new(Int32Array::from(op_values)));
390
391        RecordBatch::try_new(Arc::clone(&self.schema_with_op), columns)
392            .context("Failed to create batch with op column")
393    }
394}
395
396struct AppendEnvelopeHandler {
397    ctx: WriterContext,
398    /// Arrow schema with only user columns (no `_mz_diff`/`_mz_timestamp`), used by
399    /// [`ArrowBuilder`] to serialize row data before the extra columns are appended.
400    user_schema_for_append: Arc<ArrowSchema>,
401}
402
403impl EnvelopeHandler for AppendEnvelopeHandler {
404    fn new(
405        ctx: WriterContext,
406        _connection: &IcebergSinkConnection,
407        _materialize_arrow_schema: &Arc<ArrowSchema>,
408    ) -> anyhow::Result<Self> {
409        // arrow_schema already includes _mz_diff + _mz_timestamp (added in render_sink); strip
410        // the last two fields so ArrowBuilder only processes the user columns.
411        let n = ctx.arrow_schema.fields().len().saturating_sub(2);
412        let user_schema_for_append =
413            Arc::new(ArrowSchema::new(ctx.arrow_schema.fields()[..n].to_vec()));
414
415        Ok(Self {
416            ctx,
417            user_schema_for_append,
418        })
419    }
420
421    async fn create_writer(&self, _is_snapshot: bool) -> anyhow::Result<Box<dyn IcebergWriter>> {
422        let data_parquet_writer = ParquetWriterBuilder::new(
423            self.ctx.writer_properties.clone(),
424            Arc::clone(&self.ctx.current_schema),
425        )
426        .with_arrow_schema(Arc::clone(&self.ctx.arrow_schema))
427        .context("Arrow schema validation failed")?;
428        let data_rolling_writer = RollingFileWriterBuilder::new_with_default_file_size(
429            data_parquet_writer,
430            Arc::clone(&self.ctx.current_schema),
431            self.ctx.file_io.clone(),
432            self.ctx.location_generator.clone(),
433            self.ctx.file_name_generator.clone(),
434        );
435        Ok(Box::new(
436            DataFileWriterBuilder::new(data_rolling_writer)
437                .build(None)
438                .await
439                .context("Failed to create DataFileWriter")?,
440        ))
441    }
442
443    /// Every change is written as a plain data row: the `before` half (if present) gets
444    /// `_mz_diff = -1` and the `after` half gets `_mz_diff = +1`. Both carry the same `_mz_timestamp`.
445    fn row_to_batch(&self, diff_pair: DiffPair<Row>, ts: Timestamp) -> anyhow::Result<RecordBatch> {
446        let mut builder = ArrowBuilder::new_with_schema(
447            Arc::clone(&self.user_schema_for_append),
448            DEFAULT_ARRAY_BUILDER_ITEM_CAPACITY,
449            DEFAULT_ARRAY_BUILDER_DATA_CAPACITY,
450        )
451        .context("Failed to create builder")?;
452
453        let mut diff_values: Vec<i32> = Vec::new();
454        let ts_i64 = i64::try_from(u64::from(ts)).unwrap_or(i64::MAX);
455
456        if let Some(before) = diff_pair.before {
457            builder
458                .add_row(&before)
459                .context("Failed to add before row to builder")?;
460            diff_values.push(-1i32);
461        }
462        if let Some(after) = diff_pair.after {
463            builder
464                .add_row(&after)
465                .context("Failed to add after row to builder")?;
466            diff_values.push(1i32);
467        }
468
469        let n = diff_values.len();
470        let batch = builder
471            .to_record_batch()
472            .context("Failed to create record batch")?;
473
474        let mut columns: Vec<ArrayRef> = batch.columns().to_vec();
475        columns.push(Arc::new(Int32Array::from(diff_values)));
476        columns.push(Arc::new(Int64Array::from(vec![ts_i64; n])));
477
478        RecordBatch::try_new(Arc::clone(&self.ctx.arrow_schema), columns)
479            .context("Failed to create append record batch")
480    }
481}
482
483/// Add Parquet field IDs to an Arrow schema. Iceberg requires field IDs in the
484/// Parquet metadata for schema evolution tracking. Field IDs are assigned
485/// recursively to all nested fields (structs, lists, maps) using a depth-first,
486/// pre-order traversal.
487fn add_field_ids_to_arrow_schema(schema: ArrowSchema) -> ArrowSchema {
488    let mut next_field_id = 1i32;
489    let fields: Vec<Field> = schema
490        .fields()
491        .iter()
492        .map(|field| add_field_ids_recursive(field, &mut next_field_id))
493        .collect();
494    ArrowSchema::new(fields).with_metadata(schema.metadata().clone())
495}
496
497/// Recursively add field IDs to a field and all its nested children.
498fn add_field_ids_recursive(field: &Field, next_id: &mut i32) -> Field {
499    let current_id = *next_id;
500    *next_id += 1;
501
502    let mut metadata = field.metadata().clone();
503    metadata.insert(
504        PARQUET_FIELD_ID_META_KEY.to_string(),
505        current_id.to_string(),
506    );
507
508    let new_data_type = add_field_ids_to_datatype(field.data_type(), next_id);
509
510    Field::new(field.name(), new_data_type, field.is_nullable()).with_metadata(metadata)
511}
512
513/// Add field IDs to nested fields within a DataType.
514fn add_field_ids_to_datatype(data_type: &DataType, next_id: &mut i32) -> DataType {
515    match data_type {
516        DataType::Struct(fields) => {
517            let new_fields: Vec<Field> = fields
518                .iter()
519                .map(|f| add_field_ids_recursive(f, next_id))
520                .collect();
521            DataType::Struct(new_fields.into())
522        }
523        DataType::List(element_field) => {
524            let new_element = add_field_ids_recursive(element_field, next_id);
525            DataType::List(Arc::new(new_element))
526        }
527        DataType::LargeList(element_field) => {
528            let new_element = add_field_ids_recursive(element_field, next_id);
529            DataType::LargeList(Arc::new(new_element))
530        }
531        DataType::Map(entries_field, sorted) => {
532            let new_entries = add_field_ids_recursive(entries_field, next_id);
533            DataType::Map(Arc::new(new_entries), *sorted)
534        }
535        _ => data_type.clone(),
536    }
537}
538
539/// Merge Materialize extension metadata into Iceberg's Arrow schema.
540/// This uses Iceberg's data types (e.g. Utf8) and field IDs while preserving
541/// Materialize's extension names for ArrowBuilder compatibility.
542/// Handles nested types (structs, lists, maps) recursively.
543fn merge_materialize_metadata_into_iceberg_schema(
544    materialize_arrow_schema: &ArrowSchema,
545    iceberg_schema: &Schema,
546) -> anyhow::Result<ArrowSchema> {
547    // First, convert Iceberg schema to Arrow (this gives us the correct data types)
548    let iceberg_arrow_schema = schema_to_arrow_schema(iceberg_schema)
549        .context("Failed to convert Iceberg schema to Arrow schema")?;
550
551    // Now merge in the Materialize extension metadata
552    let fields: Vec<Field> = iceberg_arrow_schema
553        .fields()
554        .iter()
555        .map(|iceberg_field| {
556            // Find the corresponding Materialize field by name to get extension metadata
557            let mz_field = materialize_arrow_schema
558                .field_with_name(iceberg_field.name())
559                .with_context(|| {
560                    format!(
561                        "Field '{}' not found in Materialize schema",
562                        iceberg_field.name()
563                    )
564                })?;
565
566            merge_field_metadata_recursive(iceberg_field, Some(mz_field))
567        })
568        .collect::<anyhow::Result<Vec<_>>>()?;
569
570    Ok(ArrowSchema::new(fields).with_metadata(iceberg_arrow_schema.metadata().clone()))
571}
572
573/// Recursively merge Materialize extension metadata into an Iceberg field.
574fn merge_field_metadata_recursive(
575    iceberg_field: &Field,
576    mz_field: Option<&Field>,
577) -> anyhow::Result<Field> {
578    // Start with Iceberg field's metadata (which includes field IDs)
579    let mut metadata = iceberg_field.metadata().clone();
580
581    // Add Materialize extension name if available
582    if let Some(mz_f) = mz_field {
583        if let Some(extension_name) = mz_f.metadata().get(ARROW_EXTENSION_NAME_KEY) {
584            metadata.insert(ARROW_EXTENSION_NAME_KEY.to_string(), extension_name.clone());
585        }
586    }
587
588    // Recursively process nested types
589    let new_data_type = match iceberg_field.data_type() {
590        DataType::Struct(iceberg_fields) => {
591            let mz_struct_fields = match mz_field {
592                Some(f) => match f.data_type() {
593                    DataType::Struct(fields) => Some(fields),
594                    other => anyhow::bail!(
595                        "Type mismatch for field '{}': Iceberg schema has Struct, but Materialize schema has {:?}",
596                        iceberg_field.name(),
597                        other
598                    ),
599                },
600                None => None,
601            };
602
603            let new_fields: Vec<Field> = iceberg_fields
604                .iter()
605                .map(|iceberg_inner| {
606                    let mz_inner = mz_struct_fields.and_then(|fields| {
607                        fields.iter().find(|f| f.name() == iceberg_inner.name())
608                    });
609                    merge_field_metadata_recursive(iceberg_inner, mz_inner.map(|f| f.as_ref()))
610                })
611                .collect::<anyhow::Result<Vec<_>>>()?;
612
613            DataType::Struct(new_fields.into())
614        }
615        DataType::List(iceberg_element) => {
616            let mz_element = match mz_field {
617                Some(f) => match f.data_type() {
618                    DataType::List(element) => Some(element.as_ref()),
619                    other => anyhow::bail!(
620                        "Type mismatch for field '{}': Iceberg schema has List, but Materialize schema has {:?}",
621                        iceberg_field.name(),
622                        other
623                    ),
624                },
625                None => None,
626            };
627            let new_element = merge_field_metadata_recursive(iceberg_element, mz_element)?;
628            DataType::List(Arc::new(new_element))
629        }
630        DataType::LargeList(iceberg_element) => {
631            let mz_element = match mz_field {
632                Some(f) => match f.data_type() {
633                    DataType::LargeList(element) => Some(element.as_ref()),
634                    other => anyhow::bail!(
635                        "Type mismatch for field '{}': Iceberg schema has LargeList, but Materialize schema has {:?}",
636                        iceberg_field.name(),
637                        other
638                    ),
639                },
640                None => None,
641            };
642            let new_element = merge_field_metadata_recursive(iceberg_element, mz_element)?;
643            DataType::LargeList(Arc::new(new_element))
644        }
645        DataType::Map(iceberg_entries, sorted) => {
646            let mz_entries = match mz_field {
647                Some(f) => match f.data_type() {
648                    DataType::Map(entries, _) => Some(entries.as_ref()),
649                    other => anyhow::bail!(
650                        "Type mismatch for field '{}': Iceberg schema has Map, but Materialize schema has {:?}",
651                        iceberg_field.name(),
652                        other
653                    ),
654                },
655                None => None,
656            };
657            // The Iceberg arrow representation names map fields differently from
658            // Materialize (`key_value`/`key`/`value` vs `entries`/`keys`/`values`),
659            // so name-based matching on the entries struct would drop the value
660            // field's extension metadata. Merge the entries struct positionally.
661            let new_entries = match mz_entries {
662                Some(mz_entries) => merge_map_entries_metadata(iceberg_entries, mz_entries)?,
663                None => iceberg_entries.as_ref().clone(),
664            };
665            DataType::Map(Arc::new(new_entries), *sorted)
666        }
667        other => other.clone(),
668    };
669
670    Ok(Field::new(
671        iceberg_field.name(),
672        new_data_type,
673        iceberg_field.is_nullable(),
674    )
675    .with_metadata(metadata))
676}
677
678/// Merge metadata into a Map's entries struct, matching key/value positionally.
679///
680/// Iceberg's arrow representation names map fields `key_value`/`key`/`value`,
681/// while Materialize uses `entries`/`keys`/`values`. Name-based matching would
682/// drop the materialize extension metadata for the value field, which then
683/// causes `ArrowBuilder` to fail with "Field 'value' missing extension metadata".
684///
685/// Positional matching is safe because the Arrow spec defines Map structurally,
686/// not by field name: `List<entries: Struct<key: K, value: V>>` with exactly
687/// two struct children — key first, value second — and the names are only
688/// conventional. See `Map` in apache/arrow `format/Schema.fbs`:
689/// <https://github.com/apache/arrow/blob/main/format/Schema.fbs> — "The names
690/// of the child fields may be respectively 'entries', 'key', and 'value', but
691/// this is not enforced."
692///
693/// Future cleanup: we could instead align Materialize's arrow map field names
694/// with the Parquet/Iceberg convention (`key_value`/`key`/`value`) in
695/// `mz_arrow_util::builder::scalar_to_arrow_datatype_impl` and drop this
696/// positional helper. That would also affect `COPY TO S3 ... FORMAT = 'parquet'`
697/// output schemas, so we'd need to confirm no downstream consumers depend on
698/// the current `entries`/`keys`/`values` names before flipping.
699fn merge_map_entries_metadata(
700    iceberg_entries: &Field,
701    mz_entries: &Field,
702) -> anyhow::Result<Field> {
703    let mut metadata = iceberg_entries.metadata().clone();
704    if let Some(extension_name) = mz_entries.metadata().get(ARROW_EXTENSION_NAME_KEY) {
705        metadata.insert(ARROW_EXTENSION_NAME_KEY.to_string(), extension_name.clone());
706    }
707
708    let iceberg_fields = match iceberg_entries.data_type() {
709        DataType::Struct(fields) => fields,
710        other => anyhow::bail!(
711            "Iceberg map entries field '{}' is not a Struct: {:?}",
712            iceberg_entries.name(),
713            other
714        ),
715    };
716    let mz_fields = match mz_entries.data_type() {
717        DataType::Struct(fields) => fields,
718        other => anyhow::bail!(
719            "Materialize map entries field '{}' is not a Struct: {:?}",
720            mz_entries.name(),
721            other
722        ),
723    };
724
725    let new_fields: Vec<Field> = iceberg_fields
726        .iter()
727        .enumerate()
728        .map(|(idx, iceberg_inner)| {
729            let mz_inner = mz_fields.get(idx).map(|f| f.as_ref());
730            merge_field_metadata_recursive(iceberg_inner, mz_inner)
731        })
732        .collect::<anyhow::Result<Vec<_>>>()?;
733
734    Ok(Field::new(
735        iceberg_entries.name(),
736        DataType::Struct(new_fields.into()),
737        iceberg_entries.is_nullable(),
738    )
739    .with_metadata(metadata))
740}
741
742async fn reload_table(
743    catalog: &dyn Catalog,
744    namespace: String,
745    table_name: String,
746    current_table: &Table,
747) -> anyhow::Result<Table> {
748    let namespace_ident = NamespaceIdent::new(namespace.clone());
749    let table_ident = TableIdent::new(namespace_ident, table_name.clone());
750    let current_schema = current_table.metadata().current_schema_id();
751    let current_partition_spec = current_table.metadata().default_partition_spec_id();
752
753    match catalog.load_table(&table_ident).await {
754        Ok(table) => {
755            let reloaded_schema = table.metadata().current_schema_id();
756            let reloaded_partition_spec = table.metadata().default_partition_spec_id();
757            if reloaded_schema != current_schema {
758                return Err(anyhow::anyhow!(
759                    "Iceberg table '{}' schema changed during operation but schema evolution isn't supported, expected schema ID {}, got {}",
760                    table_name,
761                    current_schema,
762                    reloaded_schema
763                ));
764            }
765
766            if reloaded_partition_spec != current_partition_spec {
767                return Err(anyhow::anyhow!(
768                    "Iceberg table '{}' partition spec changed during operation but partition spec evolution isn't supported, expected partition spec ID {}, got {}",
769                    table_name,
770                    current_partition_spec,
771                    reloaded_partition_spec
772                ));
773            }
774
775            Ok(table)
776        }
777        Err(err) => Err(err).context("Failed to reload Iceberg table"),
778    }
779}
780
781/// A failed commit attempt, split by whether the commit request reached the catalog.
782enum CommitError {
783    /// Failed while building the commit locally. Nothing was sent to the catalog,
784    /// so the commit definitely didn't happen.
785    Local(iceberg::Error),
786    /// The update request failed. Depending on the error kind,
787    /// it's still possible the catalog applied the commit.
788    Request(iceberg::Error),
789}
790
791/// Build a row delta commit against the given table and send it to the catalog.
792///
793/// We can't use iceberg-rust's `Transaction::commit` (a retry wrapper for `Transaction::do_commit`)
794/// because it automatically rebases the transaction onto the latest state of the table.
795///
796/// That behavior leads to duplicate writes because there's no way to check:
797/// 1. Have we already committed this data?
798/// 2. Has another (newer) writer taken over?
799///
800/// So this implementation doesn't do that (details inline).
801async fn do_commit(
802    table: &Table,
803    catalog: &dyn Catalog,
804    snapshot_properties: Vec<(String, String)>,
805    data_files: Vec<DataFile>,
806    delete_files: Vec<DataFile>,
807) -> Result<Table, CommitError> {
808    let mut action = RowDeltaAction::new()
809        .set_snapshot_properties(snapshot_properties.into_iter().collect())
810        .with_check_duplicate(false);
811
812    if !data_files.is_empty() || !delete_files.is_empty() {
813        action = action
814            .add_data_files(data_files)
815            .add_delete_files(delete_files);
816    }
817
818    // Divergence: `Transaction::do_commit` reloads the table and rebases the transaction onto it.
819    // We do not reload the table and rebase the transaction.
820    // The caller must check the transaction against the table before committing because
821    // the Iceberg Catalog's conflict check is not sufficient to prevent duplicate commits.
822
823    let mut action_commit = Arc::new(action)
824        .commit(table)
825        .await
826        .map_err(CommitError::Local)?;
827
828    // Divergence: `Transaction::do_commit` also checks each action's requirements
829    // against the local metadata and applies its updates to a local copy of the table,
830    // so the next action in the transaction can build on the result.
831    // We only commit a single action, so we skip that step.
832
833    let table_commit = TableCommit::builder()
834        .ident(table.identifier().clone())
835        .updates(action_commit.take_updates())
836        .requirements(action_commit.take_requirements())
837        .build();
838
839    catalog
840        .update_table(table_commit)
841        .await
842        .map_err(CommitError::Request)
843}
844
845/// Attempt a single commit of a batch of data files to an Iceberg table.
846async fn try_commit_batch(
847    table: Table,
848    snapshot_properties: Vec<(String, String)>,
849    data_files: Vec<DataFile>,
850    delete_files: Vec<DataFile>,
851    catalog: &dyn Catalog,
852    conn_namespace: &str,
853    conn_table: &str,
854    sink_id: GlobalId,
855    sink_version: u64,
856    batch_lower: &Antichain<Timestamp>,
857    batch_upper: &Antichain<Timestamp>,
858    metrics: &IcebergSinkMetrics,
859) -> (Table, RetryResult<(), anyhow::Error>) {
860    // We begin the attempt by evaluating our current state.
861    // 1. Load the table from the catalog.
862    // 2. Check if the table metadata says it's safe to write to.
863    let table = match reload_table(
864        catalog,
865        conn_namespace.to_string(),
866        conn_table.to_string(),
867        &table,
868    )
869    .await
870    {
871        Ok(table) => table,
872        Err(e) => {
873            // We can't proceed without a fresh view of the table, so we must retry.
874            return (table, RetryResult::RetryableErr(anyhow!(e)));
875        }
876    };
877
878    let mut snapshots: Vec<_> = table.metadata().snapshots().cloned().collect();
879    let last = match retrieve_upper_from_snapshots(&mut snapshots) {
880        Ok(last) => last,
881        Err(e) => return (table, RetryResult::RetryableErr(anyhow!(e))),
882    };
883    if let Some((last_frontier, last_id, last_version)) = last {
884        // Just in case the sink was recreated, check both sink ID and version to see if it was us.
885        if last_id == sink_id && last_version == sink_version && last_frontier == *batch_upper {
886            // Our own commit for this batch is already on the table.
887            // We must've missed the response.
888            info!(
889                namespace = %conn_namespace,
890                table = %conn_table,
891                lower = %batch_lower.pretty(),
892                upper = %batch_upper.pretty(),
893                "found iceberg commit from previous attempt, treating as success"
894            );
895            return (table, RetryResult::Ok(()));
896        }
897
898        if last_version > sink_version {
899            // We've been superseded by a new version of the sink.
900            return (
901                table,
902                RetryResult::FatalErr(anyhow!(
903                    "Iceberg table '{}' has been modified by another writer \
904                    with version {}. Current sink version: {}. \
905                    Frontiers may be out of sync, aborting to avoid data loss.",
906                    conn_table,
907                    last_version,
908                    sink_version,
909                )),
910            );
911        }
912
913        if PartialOrder::less_equal(batch_upper, &last_frontier)
914            || PartialOrder::less_than(batch_lower, &last_frontier)
915        {
916            // This batch contains records someone else has already written.
917            return (
918                table,
919                RetryResult::FatalErr(anyhow!(
920                    "Iceberg table '{}' has been modified by another writer. \
921                    Current frontier: {:?}, last frontier: {:?}.",
922                    conn_table,
923                    batch_upper,
924                    last_frontier,
925                )),
926            );
927        }
928    }
929
930    match do_commit(
931        &table,
932        catalog,
933        snapshot_properties,
934        data_files,
935        delete_files,
936    )
937    .await
938    {
939        Ok(new_table) => (new_table, RetryResult::Ok(())),
940        Err(CommitError::Local(e)) => {
941            // Nothing was sent to the catalog, so this commit definitely didn't happen.
942            // Next attempt should reload the table and try again from scratch.
943            metrics.commit_failures.inc();
944            (
945                table,
946                RetryResult::RetryableErr(anyhow!("Failed to build iceberg table commit: {}", e)),
947            )
948        }
949        Err(CommitError::Request(e)) => match e.kind() {
950            ErrorKind::CatalogCommitConflicts => {
951                // Our view of the table was outdated.
952                // Next attempt should reload the table and try again from scratch.
953                metrics.commit_conflicts.inc();
954                (table, RetryResult::RetryableErr(anyhow!(e)))
955            }
956            ErrorKind::Unexpected => {
957                // The catalog may have applied this commit before the success response was lost.
958                // Next attempt should reload the table and see if the commit landed.
959                metrics.commit_failures.inc();
960                (table, RetryResult::RetryableErr(anyhow!(e)))
961            }
962            _ => {
963                // All other errors are definite: retrying will not change the outcome.
964                metrics.commit_failures.inc();
965                (table, RetryResult::FatalErr(anyhow!(e)))
966            }
967        },
968    }
969}
970
971/// Whether the sink can write into a table that already carries `current`,
972/// where `expected` is the schema the sink would create today.
973///
974/// Field ids, names and nullability have to match. Types have to match too,
975/// except that a column may be `fixed[16]` where `expected` has `string`. Uuid
976/// columns were written as fixed binary before, and a table holding one stays
977/// writable, because the writer builds its Arrow columns from the table's own
978/// schema rather than from `expected`. Rejecting it would strand the sink.
979///
980/// NOTE: the tolerance keys on the Iceberg types alone, so it also accepts a
981/// `fixed[16]` column where the relation has `text`. Such a table fails on the
982/// first row written instead of here. Only a table created outside Materialize
983/// can be in that state, since the sink never creates one.
984fn is_compatible(current: &Schema, expected: &Schema) -> bool {
985    current
986        .identifier_field_ids()
987        .eq(expected.identifier_field_ids())
988        && struct_is_compatible(current.as_struct(), expected.as_struct())
989}
990
991fn struct_is_compatible(current: &StructType, expected: &StructType) -> bool {
992    current.fields().len() == expected.fields().len()
993        && std::iter::zip(current.fields(), expected.fields())
994            .all(|(c, e)| field_is_compatible(c, e))
995}
996
997fn field_is_compatible(current: &NestedField, expected: &NestedField) -> bool {
998    current.id == expected.id
999        && current.name == expected.name
1000        && current.required == expected.required
1001        && match (&*current.field_type, &*expected.field_type) {
1002            // A uuid column written before uuids became strings.
1003            (Type::Primitive(PrimitiveType::Fixed(16)), Type::Primitive(PrimitiveType::String)) => {
1004                true
1005            }
1006            (Type::Struct(c), Type::Struct(e)) => struct_is_compatible(c, e),
1007            (Type::List(c), Type::List(e)) => {
1008                field_is_compatible(&c.element_field, &e.element_field)
1009            }
1010            (Type::Map(c), Type::Map(e)) => {
1011                field_is_compatible(&c.key_field, &e.key_field)
1012                    && field_is_compatible(&c.value_field, &e.value_field)
1013            }
1014            (c, e) => c == e,
1015        }
1016}
1017
1018/// Load an existing Iceberg table or create it if it doesn't exist.
1019async fn load_or_create_table(
1020    catalog: &dyn Catalog,
1021    namespace: String,
1022    table_name: String,
1023    schema: &Schema,
1024) -> anyhow::Result<iceberg::table::Table> {
1025    let namespace_ident = NamespaceIdent::new(namespace.clone());
1026    let table_ident = TableIdent::new(namespace_ident.clone(), table_name.clone());
1027
1028    // Try to load the table first
1029    match catalog.load_table(&table_ident).await {
1030        Ok(table) => {
1031            // Table exists, return it
1032            // TODO: Add proper schema evolution/validation to ensure compatibility
1033            let current_schema = table.metadata().current_schema();
1034            if !is_compatible(current_schema, schema) {
1035                anyhow::bail!(
1036                    "Iceberg table '{}' schema does not match expected schema. \
1037                     Current schema: {:?}, expected schema: {:?}",
1038                    table_name,
1039                    current_schema,
1040                    schema
1041                );
1042            }
1043            Ok(table)
1044        }
1045        Err(err) => {
1046            if matches!(err.kind(), ErrorKind::TableNotFound { .. })
1047                || err
1048                    .message()
1049                    .contains("Tried to load a table that does not exist")
1050            {
1051                // Table doesn't exist, create it
1052                // Note: location is not specified, letting the catalog determine the default location
1053                // based on its warehouse configuration
1054                let table_creation = TableCreation::builder()
1055                    .name(table_name.clone())
1056                    .schema(schema.clone())
1057                    // Use unpartitioned spec by default
1058                    // TODO: Consider making partition spec configurable
1059                    // .partition_spec(UnboundPartitionSpec::builder().build())
1060                    .build();
1061
1062                catalog
1063                    .create_table(&namespace_ident, table_creation)
1064                    .await
1065                    .with_context(|| {
1066                        format!(
1067                            "Failed to create Iceberg table '{}' in namespace '{}'",
1068                            table_name, namespace
1069                        )
1070                    })
1071            } else {
1072                // Some other error occurred
1073                Err(err).context("Failed to load Iceberg table")
1074            }
1075        }
1076    }
1077}
1078
1079/// Find the most recent Materialize frontier from Iceberg snapshots.
1080/// Returns: (frontier, sink id, sink version)
1081///
1082/// We store the frontier in snapshot metadata to track where we left off after restarts.
1083/// Snapshots with operation="replace" (compactions) don't have our metadata and are skipped.
1084/// The input slice will be sorted by sequence number in descending order.
1085fn retrieve_upper_from_snapshots(
1086    snapshots: &mut [Arc<Snapshot>],
1087) -> anyhow::Result<Option<(Antichain<Timestamp>, GlobalId, u64)>> {
1088    snapshots.sort_by(|a, b| Ord::cmp(&b.sequence_number(), &a.sequence_number()));
1089
1090    for snapshot in snapshots {
1091        let props = &snapshot.summary().additional_properties;
1092        if let (Some(frontier_json), Some(sink_id_str), Some(sink_version_str)) = (
1093            props.get("mz-frontier"),
1094            props.get("mz-sink-id"),
1095            props.get("mz-sink-version"),
1096        ) {
1097            let frontier: Vec<Timestamp> = serde_json::from_str(frontier_json)
1098                .context("Failed to deserialize frontier from snapshot properties")?;
1099            let frontier = Antichain::from_iter(frontier);
1100
1101            let sink_id = sink_id_str
1102                .parse::<GlobalId>()
1103                .context("Failed to parse mz-sink-id from snapshot properties")?;
1104
1105            let sink_version = sink_version_str
1106                .parse::<u64>()
1107                .context("Failed to parse mz-sink-version from snapshot properties")?;
1108
1109            return Ok(Some((frontier, sink_id, sink_version)));
1110        }
1111        if snapshot.summary().operation.as_str() != "replace" {
1112            // This is a bad heuristic, but we have no real other way to identify compactions
1113            // right now other than assume they will be the only operation writing "replace" operations.
1114            // That means if we find a snapshot with some other operation, but no mz-frontier, we are in an
1115            // inconsistent state and have to error out.
1116            anyhow::bail!(
1117                "Iceberg table is in an inconsistent state: snapshot {} has operation '{}' but is missing 'mz-frontier' property. Schema or partition spec evolution is not supported.",
1118                snapshot.snapshot_id(),
1119                snapshot.summary().operation.as_str(),
1120            );
1121        }
1122    }
1123
1124    Ok(None)
1125}
1126
1127/// Convert a Materialize RelationDesc into Arrow and Iceberg schemas.
1128///
1129/// Returns a tuple of:
1130/// - The Arrow schema (with field IDs and Iceberg-compatible types) for writing Parquet files
1131/// - The Iceberg schema for table creation/validation
1132///
1133/// Iceberg doesn't support unsigned integer types, so we use `iceberg_type_overrides`
1134/// to map them to compatible types (e.g., UInt64 -> Decimal128(20,0)). The ArrowBuilder
1135/// handles the cross-type conversion (Datum::UInt64 -> Decimal128Builder) automatically.
1136fn relation_desc_to_iceberg_schema(
1137    desc: &mz_repr::RelationDesc,
1138) -> anyhow::Result<(ArrowSchema, SchemaRef)> {
1139    let arrow_schema =
1140        mz_arrow_util::builder::desc_to_schema_with_overrides(desc, iceberg_type_overrides)
1141            .context("Failed to convert RelationDesc to Iceberg-compatible Arrow schema")?;
1142
1143    let arrow_schema_with_ids = add_field_ids_to_arrow_schema(arrow_schema);
1144
1145    let iceberg_schema = arrow_schema_to_schema(&arrow_schema_with_ids)
1146        .context("Failed to convert Arrow schema to Iceberg schema")?;
1147
1148    Ok((arrow_schema_with_ids, Arc::new(iceberg_schema)))
1149}
1150
1151/// Resolve Materialize key column indexes to Iceberg top-level field IDs.
1152///
1153/// Iceberg field IDs are assigned recursively, so a top-level column's field ID
1154/// is not necessarily `column_index + 1` once nested fields are present.
1155fn equality_ids_for_indices(
1156    current_schema: &Schema,
1157    materialize_arrow_schema: &ArrowSchema,
1158    equality_indices: &[usize],
1159) -> anyhow::Result<Vec<i32>> {
1160    let top_level_fields = current_schema.as_struct();
1161
1162    equality_indices
1163        .iter()
1164        .map(|index| {
1165            let mz_field = materialize_arrow_schema
1166                .fields()
1167                .get(*index)
1168                .with_context(|| format!("Equality delete key index {index} is out of bounds"))?;
1169            let field_name = mz_field.name();
1170            let iceberg_field = top_level_fields
1171                .field_by_name(field_name)
1172                .with_context(|| {
1173                    format!(
1174                        "Equality delete key column '{}' not found in Iceberg table schema",
1175                        field_name
1176                    )
1177                })?;
1178            Ok(iceberg_field.id)
1179        })
1180        .collect()
1181}
1182
1183/// Build a new Arrow schema by adding an __op column to the existing schema.
1184fn build_schema_with_op_column(schema: &ArrowSchema) -> ArrowSchema {
1185    let mut fields: Vec<Arc<Field>> = schema.fields().iter().cloned().collect();
1186    fields.push(Arc::new(Field::new("__op", DataType::Int32, false)));
1187    ArrowSchema::new(fields)
1188}
1189
1190/// Build a new Arrow schema by appending `_mz_diff` (Int32) and `_mz_timestamp` (Int64) columns.
1191/// These are user-visible Iceberg columns written in append mode. Parquet field IDs are
1192/// assigned sequentially after the existing maximum field ID so the extended schema can
1193/// be converted to a valid Iceberg schema via `arrow_schema_to_schema`.
1194#[allow(clippy::disallowed_types)]
1195fn build_schema_with_append_columns(schema: &ArrowSchema) -> ArrowSchema {
1196    use mz_storage_types::sinks::{ICEBERG_APPEND_DIFF_COLUMN, ICEBERG_APPEND_TIMESTAMP_COLUMN};
1197    let mut fields: Vec<Arc<Field>> = schema.fields().iter().cloned().collect();
1198    fields.push(Arc::new(Field::new(
1199        ICEBERG_APPEND_DIFF_COLUMN,
1200        DataType::Int32,
1201        false,
1202    )));
1203    fields.push(Arc::new(Field::new(
1204        ICEBERG_APPEND_TIMESTAMP_COLUMN,
1205        DataType::Int64,
1206        false,
1207    )));
1208
1209    add_field_ids_to_arrow_schema(ArrowSchema::new(fields).with_metadata(schema.metadata().clone()))
1210}
1211
1212/// Generate time-based batch boundaries for grouping writes into Iceberg snapshots.
1213/// Batches are minted with configurable windows to balance write efficiency with latency.
1214/// We maintain a sliding window of future batch descriptions so writers can start
1215/// processing data even while earlier batches are still being written.
1216fn mint_batch_descriptions<'scope>(
1217    name: String,
1218    sink_id: GlobalId,
1219    input: SinkBatchStream<'scope>,
1220    sink: &StorageSinkDesc<CollectionMetadata, Timestamp>,
1221    connection: IcebergSinkConnection,
1222    storage_configuration: StorageConfiguration,
1223    initial_schema: SchemaRef,
1224) -> (
1225    StreamVec<'scope, Timestamp, (Antichain<Timestamp>, Antichain<Timestamp>)>,
1226    StreamVec<'scope, Timestamp, Infallible>,
1227    StreamVec<'scope, Timestamp, HealthStatusMessage>,
1228    PressOnDropButton,
1229) {
1230    let scope = input.scope();
1231    let name_for_error = name.clone();
1232    let name_for_logging = name.clone();
1233    let mut builder = OperatorBuilder::new(name, scope.clone());
1234    let sink_version = sink.version;
1235
1236    let hashed_id = sink_id.hashed();
1237    let is_active_worker = usize::cast_from(hashed_id) % scope.peers() == scope.index();
1238    let (_, table_ready_stream) = builder.new_output::<CapacityContainerBuilder<Vec<_>>>();
1239    let (batch_desc_output, batch_desc_stream) =
1240        builder.new_output::<CapacityContainerBuilder<Vec<_>>>();
1241    let mut input = builder.new_input_for(input, Pipeline, &batch_desc_output);
1242
1243    let as_of = sink.as_of.clone();
1244    let commit_interval = sink
1245        .commit_interval
1246        .expect("the planner should have enforced this")
1247        .clone();
1248
1249    let (button, errors): (_, StreamVec<'scope, Timestamp, Rc<anyhow::Error>>) =
1250        builder.build_fallible(move |caps| {
1251        Box::pin(async move {
1252            let [table_ready_capset, capset]: &mut [_; 2] = caps.try_into().unwrap();
1253
1254            if !is_active_worker {
1255                // Only the active worker mints batch descriptions.
1256                return Ok(());
1257            }
1258
1259            let table_ident = TableIdent::new(
1260                NamespaceIdent::new(connection.namespace.clone()),
1261                connection.table.clone(),
1262            );
1263            let catalog = connection
1264                .catalog_connection
1265                .connect(&storage_configuration, InTask::Yes, Some(&table_ident))
1266                .await
1267                .with_context(|| {
1268                    format!(
1269                        "Failed to connect to Iceberg catalog '{}' for table '{}.{}'",
1270                        connection.catalog_connection.uri, connection.namespace, connection.table
1271                    )
1272                })?;
1273
1274            let table = load_or_create_table(
1275                catalog.as_ref(),
1276                connection.namespace.clone(),
1277                connection.table.clone(),
1278                initial_schema.as_ref(),
1279            )
1280            .await?;
1281            debug!(
1282                ?sink_id,
1283                %name_for_logging,
1284                namespace = %connection.namespace,
1285                table = %connection.table,
1286                "iceberg mint loaded/created table"
1287            );
1288
1289            *table_ready_capset = CapabilitySet::new();
1290
1291            let mut snapshots: Vec<_> = table.metadata().snapshots().cloned().collect();
1292            let resume = retrieve_upper_from_snapshots(&mut snapshots)?;
1293            let (resume_upper, resume_version) = match resume {
1294                Some((f, _, v)) => (f, v),
1295                None => (Antichain::from_elem(Timestamp::minimum()), 0),
1296            };
1297            debug!(
1298                ?sink_id,
1299                %name_for_logging,
1300                resume_upper = %resume_upper.pretty(),
1301                resume_version,
1302                as_of = %as_of.pretty(),
1303                "iceberg mint resume position loaded"
1304            );
1305
1306            // The input has overcompacted if
1307            let overcompacted =
1308                // ..we have made some progress in the past
1309                *resume_upper != [Timestamp::minimum()] &&
1310                // ..but the since frontier is now beyond that
1311                PartialOrder::less_than(&resume_upper, &as_of);
1312
1313            if overcompacted {
1314                let err = format!(
1315                    "{name_for_error}: input compacted past resume upper: as_of {}, resume_upper: {}",
1316                    as_of.pretty(),
1317                    resume_upper.pretty()
1318                );
1319                // This would normally be an assertion but because it can happen after a
1320                // Materialize backup/restore we log an error so that it appears on Sentry but
1321                // leaves the rest of the objects in the cluster unaffected.
1322                return Err(anyhow::anyhow!("{err}"));
1323            };
1324
1325            if resume_version > sink_version {
1326                anyhow::bail!("Fenced off by newer sink version: resume_version {}, sink_version {}", resume_version, sink_version);
1327            }
1328
1329            let mut initialized = false;
1330            let mut observed_frontier;
1331            // Track minted batches to maintain a sliding window of open batch descriptions.
1332            // This is needed to know when to retire old batches and mint new ones.
1333            // It's "sortedness" is derived from the monotonicity of batch descriptions,
1334            // and the fact that we only ever push new descriptions to the back and pop from the front.
1335            let mut minted_batches = VecDeque::new();
1336
1337            // Once we start seeing new data, we'll roll everything into a single catch-up commit
1338            // before beginning the steady state commit interval.
1339            let catchup_start = if *resume_upper == [Timestamp::minimum()] {
1340                // If we're hydrating from a source snapshot, we immediately emit the snapshot's batch description.
1341                let batch_upper = Antichain::from_elem(
1342                    as_of.as_option().expect("as_of not empty").step_forward());
1343                let batch = (as_of.clone(), batch_upper.clone());
1344                minted_batches.push_back(batch.clone());
1345                batch_desc_output.give(&capset[0], batch);
1346                capset.downgrade(batch_upper.clone());
1347
1348                // The "catch-up" batch starts at the end of the snapshot batch.
1349                batch_upper
1350            } else {
1351                // If we're resuming, the "catch-up" batch starts at the start of data, i.e. resume_upper.
1352                resume_upper.clone()
1353            };
1354
1355            loop {
1356                if let Some(event) = input.next().await {
1357                    match event {
1358                        Event::Data(_, _) => continue,
1359                        Event::Progress(frontier) => {
1360                            observed_frontier = frontier;
1361                        }
1362                    }
1363                } else {
1364                    return Ok(());
1365                }
1366
1367                if !initialized {
1368                    if observed_frontier.is_empty() {
1369                        // Bounded inputs can close (frontier becomes empty) before we finish
1370                        // initialization. For example, a loadgen source configured for a finite
1371                        // dataset may emit all rows at time t and then immediately close.
1372                        // Mint one final batch with an empty upper. The input is closed, so
1373                        // that batch covers all remaining data on every worker, and committing
1374                        // it records the sink as complete.
1375                        if catchup_start.is_empty() {
1376                            // A previous incarnation already committed through the empty
1377                            // frontier. Nothing left to do.
1378                            return Ok(());
1379                        }
1380                        debug!(
1381                            ?sink_id,
1382                            %name_for_logging,
1383                            batch_lower = %catchup_start.pretty(),
1384                            "iceberg mint input closed before initialization; minting final batch"
1385                        );
1386                        let batch = (catchup_start.clone(), Antichain::new());
1387                        batch_desc_output.give(&capset[0], batch);
1388                        return Ok(());
1389                    }
1390
1391                    // Don't make empty commits while we wait ^for the first data to be ready.
1392                    // (^for the frontier to indicate there _could_ be data ready)
1393                    if !PartialOrder::less_than(&catchup_start, &observed_frontier)
1394                    {
1395                        continue;
1396                    }
1397
1398                    let mut batch_descriptions = vec![];
1399                    let mut current_upper = observed_frontier.clone();
1400                    let current_upper_ts = observed_frontier.as_option().expect("frontier not empty").clone();
1401                    debug!(
1402                        ?sink_id,
1403                        %name_for_logging,
1404                        batch_lower = %catchup_start.pretty(),
1405                        current_upper = %current_upper.pretty(),
1406                        "iceberg mint initializing (catch-up batch)"
1407                    );
1408                    debug!(
1409                        "{}: creating catch-up batch [{}, {})",
1410                        name_for_logging,
1411                        catchup_start.pretty(),
1412                        current_upper.pretty()
1413                    );
1414                    batch_descriptions.push((catchup_start.clone(), current_upper.clone()));
1415
1416                    // Mint initial future batch descriptions at configurable intervals
1417                    for i in 1..INITIAL_DESCRIPTIONS_TO_MINT {
1418                        let duration_millis = commit_interval.as_millis()
1419                            .checked_mul(u128::from(i))
1420                            .expect("commit interval multiplication overflow");
1421                        let duration_ts = Timestamp::new(
1422                            u64::try_from(duration_millis)
1423                                .expect("commit interval too large for u64"),
1424                        );
1425                        let desired_batch_upper = Antichain::from_elem(
1426                            current_upper_ts.step_forward_by(&duration_ts),
1427                        );
1428
1429                        let batch_description =
1430                            (current_upper.clone(), desired_batch_upper.clone());
1431                        debug!(
1432                            "{}: minting future batch {}/{} [{}, {})",
1433                            name_for_logging,
1434                            i,
1435                            INITIAL_DESCRIPTIONS_TO_MINT,
1436                            current_upper.pretty(),
1437                            desired_batch_upper.pretty()
1438                        );
1439                        current_upper = batch_description.1.clone();
1440                        batch_descriptions.push(batch_description);
1441                    }
1442
1443                    minted_batches.extend(batch_descriptions.clone());
1444
1445                    for desc in batch_descriptions {
1446                        batch_desc_output.give(&capset[0], desc);
1447                    }
1448
1449                    capset.downgrade(current_upper);
1450
1451                    initialized = true;
1452                } else {
1453                    if observed_frontier.is_empty() {
1454                        // We're done!
1455                        return Ok(());
1456                    }
1457                    // Maintain a sliding window: when the oldest batch becomes ready, retire it
1458                    // and mint a new future batch to keep the pipeline full
1459                    while let Some(oldest_desc) = minted_batches.front() {
1460                        let oldest_upper = &oldest_desc.1;
1461                        if !PartialOrder::less_equal(oldest_upper, &observed_frontier) {
1462                            break;
1463                        }
1464
1465                        let newest_upper = minted_batches.back().unwrap().1.clone();
1466                        let new_lower = newest_upper.clone();
1467                        let duration_ts = Timestamp::new(commit_interval.as_millis()
1468                            .try_into()
1469                            .expect("commit interval too large for u64"));
1470                        let new_upper = Antichain::from_elem(newest_upper
1471                            .as_option()
1472                            .unwrap()
1473                            .step_forward_by(&duration_ts));
1474
1475                        let new_batch_description = (new_lower.clone(), new_upper.clone());
1476                        minted_batches.pop_front();
1477                        minted_batches.push_back(new_batch_description.clone());
1478
1479                        batch_desc_output.give(&capset[0], new_batch_description);
1480
1481                        capset.downgrade(new_upper);
1482                    }
1483                }
1484            }
1485        })
1486    });
1487
1488    let statuses = errors.map(|error| HealthStatusMessage {
1489        id: None,
1490        update: HealthStatusUpdate::halting(format!("{}", error.display_with_causes()), None),
1491        namespace: StatusNamespace::Iceberg,
1492    });
1493    (
1494        batch_desc_stream,
1495        table_ready_stream,
1496        statuses,
1497        button.press_on_drop(),
1498    )
1499}
1500
1501#[derive(Clone, Debug, Serialize, Deserialize)]
1502#[serde(try_from = "AvroDataFile", into = "AvroDataFile")]
1503struct SerializableDataFile {
1504    pub data_file: DataFile,
1505    pub schema: Schema,
1506}
1507
1508/// A wrapper around Iceberg's DataFile that implements Serialize and Deserialize.
1509/// This is slightly complicated by the fact that Iceberg's DataFile doesn't implement
1510/// these traits directly, so we serialize to/from Avro bytes (which Iceberg supports natively).
1511/// The avro ser(de) also requires the Iceberg schema to be provided, so we include that as well.
1512/// It is distinctly possible that this is overkill, but it avoids re-implementing
1513/// Iceberg's serialization logic here.
1514/// If at some point this becomes a serious overhead, we can revisit this decision.
1515#[derive(Clone, Debug, Serialize, Deserialize)]
1516struct AvroDataFile {
1517    pub data_file: Vec<u8>,
1518    /// Schema serialized as JSON bytes to avoid bincode issues with HashMap
1519    pub schema: Vec<u8>,
1520}
1521
1522impl From<SerializableDataFile> for AvroDataFile {
1523    fn from(value: SerializableDataFile) -> Self {
1524        let mut data_file = Vec::new();
1525        write_data_files_to_avro(
1526            &mut data_file,
1527            [value.data_file],
1528            &StructType::new(vec![]),
1529            FormatVersion::V2,
1530        )
1531        .expect("serialization into buffer");
1532        let schema = serde_json::to_vec(&value.schema).expect("schema serialization");
1533        AvroDataFile { data_file, schema }
1534    }
1535}
1536
1537impl TryFrom<AvroDataFile> for SerializableDataFile {
1538    type Error = String;
1539
1540    fn try_from(value: AvroDataFile) -> Result<Self, Self::Error> {
1541        let schema: Schema = serde_json::from_slice(&value.schema)
1542            .map_err(|e| format!("Failed to deserialize schema: {}", e))?;
1543        let data_files = read_data_files_from_avro(
1544            &mut &*value.data_file,
1545            &schema,
1546            0,
1547            &StructType::new(vec![]),
1548            FormatVersion::V2,
1549        )
1550        .map_err_to_string_with_causes()?;
1551        let Some(data_file) = data_files.into_iter().next() else {
1552            return Err("No DataFile found in Avro data".into());
1553        };
1554        Ok(SerializableDataFile { data_file, schema })
1555    }
1556}
1557
1558/// A DataFile along with its associated batch description (lower and upper bounds).
1559#[derive(Clone, Debug, Serialize, Deserialize)]
1560struct BoundedDataFile {
1561    pub data_file: SerializableDataFile,
1562    pub batch_desc: (Antichain<Timestamp>, Antichain<Timestamp>),
1563}
1564
1565impl BoundedDataFile {
1566    pub fn new(
1567        file: DataFile,
1568        schema: Schema,
1569        batch_desc: (Antichain<Timestamp>, Antichain<Timestamp>),
1570    ) -> Self {
1571        Self {
1572            data_file: SerializableDataFile {
1573                data_file: file,
1574                schema,
1575            },
1576            batch_desc,
1577        }
1578    }
1579
1580    pub fn batch_desc(&self) -> &(Antichain<Timestamp>, Antichain<Timestamp>) {
1581        &self.batch_desc
1582    }
1583
1584    pub fn data_file(&self) -> &DataFile {
1585        &self.data_file.data_file
1586    }
1587
1588    pub fn into_data_file(self) -> DataFile {
1589        self.data_file.data_file
1590    }
1591}
1592
1593/// A set of DataFiles along with their associated batch descriptions.
1594#[derive(Clone, Debug, Default)]
1595struct BoundedDataFileSet {
1596    pub data_files: Vec<BoundedDataFile>,
1597}
1598
1599/// Returns the base location for the table's data files, with no trailing separator.
1600///
1601/// `configured_path` is the catalog's `write.data.path`, or `write.folder-storage.path` where
1602/// only the older property is set. `location` is the table's own location, used when the catalog
1603/// configures neither.
1604///
1605/// The result never ends in `/`. Callers join it with a `/` and a file name, and the joined URI
1606/// is what lands in the manifest, so a separator left on the end here produces a manifest entry
1607/// naming an object that was never written.
1608fn data_file_location(configured_path: Option<&str>, location: &str) -> String {
1609    // Both properties may legally end in `/`. `DefaultLocationGenerator` stores the value
1610    // verbatim and `generate_location` appends `/` plus the file name, so `s3://b/t/data/`
1611    // yields `s3://b/t/data//f.parquet`. OpenDAL collapses the `//` when it writes the object,
1612    // but the Parquet writer copies the unnormalized URI into the `DataFile`, leaving the
1613    // manifest pointing at a key that does not exist and the table unreadable to anyone else.
1614    // The reference Iceberg location provider strips them for this reason.
1615    if let Some(path) = configured_path {
1616        return path.trim_end_matches('/').to_string();
1617    }
1618
1619    // WORKAROUND: S3 Tables catalog incorrectly sets location to the metadata file path
1620    // instead of the warehouse root. Strip off the /metadata/*.metadata.json suffix. No
1621    // clear way to detect this properly right now, so we use heuristics.
1622    let corrected_location = match location.rsplit_once("/metadata/") {
1623        Some((a, b)) if b.ends_with(".metadata.json") => a,
1624        _ => location,
1625    };
1626    // Trimmed before the join, not after, or a location ending in `/` moves the doubled
1627    // separator into the middle of the URI where a trailing trim cannot reach it.
1628    format!("{}/data", corrected_location.trim_end_matches('/'))
1629}
1630
1631/// Construct the envelope-specific closures that [`write_data_files`] needs.
1632///
1633/// Write rows into Parquet data files bounded by batch descriptions.
1634/// Rows are matched to batches by timestamp; if a batch description hasn't arrived yet,
1635/// rows are stashed until it does. This allows batches to be minted ahead of data arrival.
1636fn write_data_files<'scope, H: EnvelopeHandler + 'static>(
1637    name: String,
1638    input: SinkBatchStream<'scope>,
1639    batch_desc_input: StreamVec<'scope, Timestamp, (Antichain<Timestamp>, Antichain<Timestamp>)>,
1640    table_ready_stream: StreamVec<'scope, Timestamp, Infallible>,
1641    sink_id: GlobalId,
1642    from_id: GlobalId,
1643    key_is_synthetic: bool,
1644    as_of: Antichain<Timestamp>,
1645    connection: IcebergSinkConnection,
1646    storage_configuration: StorageConfiguration,
1647    materialize_arrow_schema: Arc<ArrowSchema>,
1648    metrics: Arc<IcebergSinkMetrics>,
1649    statistics: SinkStatistics,
1650) -> (
1651    StreamVec<'scope, Timestamp, BoundedDataFile>,
1652    StreamVec<'scope, Timestamp, HealthStatusMessage>,
1653    PressOnDropButton,
1654) {
1655    let scope = input.scope();
1656    let name_for_logging = name.clone();
1657    let mut builder = OperatorBuilder::new(name, scope.clone());
1658
1659    let (output, output_stream) = builder.new_output::<CapacityContainerBuilder<_>>();
1660
1661    let mut table_ready_input = builder.new_disconnected_input(table_ready_stream, Pipeline);
1662    let mut batch_desc_input =
1663        builder.new_input_for(batch_desc_input.broadcast(), Pipeline, &output);
1664    let mut input = builder.new_disconnected_input(input, Pipeline);
1665
1666    let (button, errors): (_, StreamVec<'scope, Timestamp, Rc<anyhow::Error>>) = builder
1667        .build_fallible(move |caps| {
1668            Box::pin(async move {
1669                let [capset]: &mut [_; 1] = caps.try_into().unwrap();
1670                let namespace_ident = NamespaceIdent::new(connection.namespace.clone());
1671                let table_ident = TableIdent::new(namespace_ident, connection.table.clone());
1672                let catalog = connection
1673                    .catalog_connection
1674                    .connect(&storage_configuration, InTask::Yes, Some(&table_ident))
1675                    .await
1676                    .with_context(|| {
1677                        format!(
1678                            "Failed to connect to Iceberg catalog '{}' for table '{}.{}'",
1679                            connection.catalog_connection.uri,
1680                            connection.namespace,
1681                            connection.table
1682                        )
1683                    })?;
1684
1685                while let Some(_) = table_ready_input.next().await {
1686                    // Wait for table to be ready
1687                }
1688                let table = catalog.load_table(&table_ident).await.with_context(|| {
1689                    format!(
1690                        "Failed to load Iceberg table '{}.{}' in write_data_files operator",
1691                        connection.namespace, connection.table
1692                    )
1693                })?;
1694
1695                let table_metadata = table.metadata().clone();
1696                let current_schema = Arc::clone(table_metadata.current_schema());
1697
1698                // Merge Materialize extension metadata into the Iceberg schema.
1699                // We need extension metadata for ArrowBuilder to work correctly (it uses
1700                // extension names to know how to handle different types like records vs arrays).
1701                let arrow_schema = Arc::new(
1702                    merge_materialize_metadata_into_iceberg_schema(
1703                        materialize_arrow_schema.as_ref(),
1704                        current_schema.as_ref(),
1705                    )
1706                    .context("Failed to merge Materialize metadata into Iceberg schema")?,
1707                );
1708
1709                // A catalog that manages where data files live advertises it through
1710                // `write.data.path`. Honor it: catalogs backing an Iceberg table with
1711                // their own storage layout reject a commit whose data files sit outside
1712                // that path. Unity Catalog, for one, has to register the files in the
1713                // Delta log that actually backs the table, and answers a commit
1714                // referencing files under `<location>/data` with a 500.
1715                //
1716                // `DefaultLocationGenerator::new` reads these same properties, but its
1717                // fallback misses the S3 Tables correction, so choose explicitly.
1718                let properties = table_metadata.properties();
1719                let configured_path = properties
1720                    .get("write.data.path")
1721                    .or_else(|| properties.get("write.folder-storage.path"));
1722                let data_location = data_file_location(
1723                    configured_path.map(String::as_str),
1724                    table_metadata.location(),
1725                );
1726                debug!(%data_location, "iceberg sink data file location");
1727                let location_generator =
1728                    DefaultLocationGenerator::with_data_location(data_location);
1729
1730                // Add a unique suffix to avoid filename collisions across restarts and workers
1731                let unique_suffix = format!("-{}", uuid::Uuid::new_v4());
1732                let file_name_generator = DefaultFileNameGenerator::new(
1733                    PARQUET_FILE_PREFIX.to_string(),
1734                    Some(unique_suffix),
1735                    iceberg::spec::DataFileFormat::Parquet,
1736                );
1737
1738                let file_io = table.file_io().clone();
1739
1740                let writer_properties = WriterProperties::new();
1741
1742                let ctx = WriterContext {
1743                    arrow_schema,
1744                    current_schema: Arc::clone(&current_schema),
1745                    file_io,
1746                    location_generator,
1747                    file_name_generator,
1748                    writer_properties,
1749                };
1750                let handler = H::new(ctx, &connection, &materialize_arrow_schema)?;
1751                let mut pk_warner =
1752                    (!key_is_synthetic).then(|| PkViolationWarner::new(sink_id, from_id));
1753
1754                // Rows can arrive before their batch description due to dataflow parallelism.
1755                // Stash them until we know which batch they belong to.
1756                // Keyed by the lower bound (per arrangement batch) of the rows.
1757                let mut stashed_rows: VecDeque<ArcBatch<OrdValBatch<_>>> = VecDeque::new();
1758
1759                // Track batches currently being written. When a row arrives, we check if it belongs
1760                // to an in-flight batch. When frontiers advance to a batch's upper, we close the
1761                // writer and emit its data files downstream.
1762                let mut in_flight_batches: VecDeque<(
1763                    (Antichain<Timestamp>, Antichain<Timestamp>),
1764                    Box<dyn IcebergWriter>,
1765                )> = VecDeque::new();
1766
1767                // The bounds of the most recently received batch description and input batch.
1768                // `with_ready_batches` relies on both inputs arriving in order and
1769                // non-overlapping. These track that invariant for the checks below.
1770                let mut last_batch_desc: Option<BatchDescription> = None;
1771                let mut last_input_bounds: Option<(Antichain<Timestamp>, Antichain<Timestamp>)> =
1772                    None;
1773
1774                let mut batch_description_frontier = Antichain::from_elem(Timestamp::minimum());
1775                let mut input_frontier = Antichain::from_elem(Timestamp::minimum());
1776
1777                while !(batch_description_frontier.is_empty() && input_frontier.is_empty()) {
1778                    tokio::select! {
1779                        _ = batch_desc_input.ready() => {},
1780                        _ = input.ready() => {}
1781                    }
1782
1783                    // Operator Recipe Step 1: Read all the input.
1784
1785                    // Read all the incoming batch descriptions.
1786                    while let Some(event) = batch_desc_input.next_sync() {
1787                        match event {
1788                            Event::Data(_cap, data) => {
1789                                for batch_desc in data {
1790                                    let (lower, upper) = &batch_desc;
1791
1792                                    if let Some((prev_lower, prev_upper)) = last_batch_desc.as_ref()
1793                                    {
1794                                        if prev_upper != lower {
1795                                            anyhow::bail!(
1796                                                "batch descriptions must arrive in order, non-overlapping, \
1797                                                and without gaps: previous [{}, {}), new [{}, {})",
1798                                                prev_lower.pretty(),
1799                                                prev_upper.pretty(),
1800                                                lower.pretty(),
1801                                                upper.pretty(),
1802                                            );
1803                                        }
1804                                    }
1805                                    last_batch_desc = Some(batch_desc.clone());
1806
1807                                    // Disable seen_rows tracking for snapshot batch to save memory
1808                                    let is_snapshot = lower == &as_of;
1809                                    debug!(
1810                                        "{}: received batch description [{}, {}), snapshot={}",
1811                                        name_for_logging,
1812                                        lower.pretty(),
1813                                        upper.pretty(),
1814                                        is_snapshot
1815                                    );
1816                                    let batch_writer = handler.create_writer(is_snapshot).await?;
1817                                    in_flight_batches.push_back((batch_desc.clone(), batch_writer));
1818                                }
1819                            }
1820                            Event::Progress(frontier) => {
1821                                batch_description_frontier = frontier;
1822                            }
1823                        }
1824                    }
1825
1826                    // Read all the incoming (arrangement batches of) rows.
1827                    while let Some(event) = input.next_sync() {
1828                        match event {
1829                            Event::Data(_cap, data) => {
1830                                for rows in &data {
1831                                    if let Some((prev_lower, prev_upper)) =
1832                                        last_input_bounds.as_ref()
1833                                    {
1834                                        // If the collection doesn't change and the frontier advances,
1835                                        // we can (correctly) observe gaps between input batches.
1836                                        // This differs from output batch descriptions, which are constructed without gaps.
1837                                        if !PartialOrder::less_equal(prev_upper, rows.lower()) {
1838                                            anyhow::bail!(
1839                                                "input batches must arrive in order and \
1840                                                non-overlapping: previous [{}, {}), new [{}, {})",
1841                                                prev_lower.pretty(),
1842                                                prev_upper.pretty(),
1843                                                rows.lower().pretty(),
1844                                                rows.upper().pretty(),
1845                                            );
1846                                        }
1847                                    }
1848                                    last_input_bounds =
1849                                        Some((rows.lower().clone(), rows.upper().clone()));
1850
1851                                    stashed_rows.push_back(rows.clone());
1852                                }
1853                            }
1854                            Event::Progress(frontier) => {
1855                                input_frontier = frontier;
1856                            }
1857                        }
1858                    }
1859
1860                    metrics.stashed_rows.set(u64::cast_from(
1861                        stashed_rows.iter().map(|rows| rows.len()).sum::<usize>(),
1862                    ));
1863
1864                    // Operator Recipe Steps 2-4: Consult frontiers. Plan work. Do all the work.
1865
1866                    // Report staged messages periodically during writes so progress is
1867                    // visible while a large batch is still open.
1868                    let mut staged_messages_since_flush: u64 = 0;
1869
1870                    // How to write rows from a(n arrangement) batch into a(n Iceberg) batch.
1871                    let write_rows = async |rows: &OrdValBatch<_>,
1872                                            (lower, upper): BatchDescription,
1873                                            batch_writer: &mut Box<dyn IcebergWriter>|
1874                           -> Result<(), anyhow::Error> {
1875                        for_each_diff_pair_async(
1876                            rows,
1877                            Some(lower),
1878                            Some(upper),
1879                            async |key, time, diff_pair| -> Result<(), anyhow::Error> {
1880                                if let Some(warner) = pk_warner.as_mut() {
1881                                    warner.observe(key, time);
1882                                }
1883
1884                                let record_batch = handler
1885                                    .row_to_batch(diff_pair, time)
1886                                    .context("failed to convert row to recordbatch")?;
1887                                staged_messages_since_flush +=
1888                                    u64::cast_from(record_batch.num_rows());
1889                                batch_writer
1890                                    .write(record_batch)
1891                                    .await
1892                                    .context("failed to write recordbatch")?;
1893                                if staged_messages_since_flush >= 10_000 {
1894                                    statistics.inc_messages_staged_by(staged_messages_since_flush);
1895                                    staged_messages_since_flush = 0;
1896                                }
1897                                Ok(())
1898                            },
1899                        )
1900                        .await?;
1901                        // Flush after each batch so the final `(key, time)` group of the walk is
1902                        // resolved immediately — a PK violation in the last group is otherwise held
1903                        // until more data arrives or the operator shuts down.
1904                        if let Some(warner) = pk_warner.as_mut() {
1905                            warner.flush();
1906                        }
1907                        Ok(())
1908                    };
1909
1910                    // How to seal the data files for an Iceberg commit.
1911                    let close_batch = async |batch_desc: BatchDescription,
1912                                             batch_writer: &mut Box<dyn IcebergWriter>|
1913                           -> Result<(), anyhow::Error> {
1914                        let close_started_at = Instant::now();
1915                        let data_files = batch_writer.close().await;
1916                        metrics
1917                            .writer_close_duration_seconds
1918                            .observe(close_started_at.elapsed().as_secs_f64());
1919                        let data_files = data_files.context("Failed to close batch writer")?;
1920                        debug!(
1921                            "{}: closed batch [{}, {}), wrote {} files",
1922                            name_for_logging,
1923                            batch_desc.0.pretty(),
1924                            batch_desc.1.pretty(),
1925                            data_files.len()
1926                        );
1927                        for data_file in data_files {
1928                            match data_file.content_type() {
1929                                iceberg::spec::DataContentType::Data => {
1930                                    metrics.data_files_written.inc();
1931                                }
1932                                iceberg::spec::DataContentType::PositionDeletes
1933                                | iceberg::spec::DataContentType::EqualityDeletes => {
1934                                    metrics.delete_files_written.inc();
1935                                }
1936                            }
1937                            statistics.inc_bytes_staged_by(data_file.file_size_in_bytes());
1938                            let file = BoundedDataFile::new(
1939                                data_file,
1940                                current_schema.as_ref().clone(),
1941                                batch_desc.clone(),
1942                            );
1943                            output.give(&capset[0], file);
1944                        }
1945
1946                        // Operator Recipe Step 5: Downgrade or drop capabilities.
1947
1948                        capset.downgrade(batch_desc.1.clone());
1949                        Ok(())
1950                    };
1951
1952                    // Write the rows and seal the data files.
1953                    with_ready_batches(
1954                        input_frontier.clone(),
1955                        &mut stashed_rows,
1956                        batch_description_frontier.clone(),
1957                        &mut in_flight_batches,
1958                        write_rows,
1959                        close_batch,
1960                    )
1961                    .await?;
1962
1963                    if staged_messages_since_flush > 0 {
1964                        statistics.inc_messages_staged_by(staged_messages_since_flush);
1965                    }
1966                    metrics.stashed_rows.set(u64::cast_from(
1967                        stashed_rows.iter().map(|rows| rows.len()).sum::<usize>(),
1968                    ));
1969                }
1970                Ok(())
1971            })
1972        });
1973
1974    let statuses = errors.map(|error| HealthStatusMessage {
1975        id: None,
1976        update: HealthStatusUpdate::halting(format!("{}", error.display_with_causes()), None),
1977        namespace: StatusNamespace::Iceberg,
1978    });
1979    (output_stream, statuses, button.press_on_drop())
1980}
1981
1982/// The `[lower, upper)` frontier bounds of one Iceberg commit.
1983type BatchDescription = (Antichain<Timestamp>, Antichain<Timestamp>);
1984
1985/// Write out as much of the input as we can.
1986///
1987/// Drop input batches when:
1988/// - their contents have all been written out
1989/// - no possible future output batch could need their contents
1990///
1991/// Close and drop output batches when:
1992/// - no possible future input batch could overlap with their time window
1993///
1994/// Invariant: We assume the batches in each stream (input vs output)
1995/// are in order and non-overlapping.
1996async fn with_ready_batches<L: Layout, W, Write, Close>(
1997    input_frontier: Antichain<Timestamp>,
1998    input_batches: &mut VecDeque<ArcBatch<OrdValBatch<L>>>,
1999    output_frontier: Antichain<Timestamp>,
2000    output_batches: &mut VecDeque<(BatchDescription, W)>,
2001    mut write_rows: Write,
2002    mut close_batch: Close,
2003) -> Result<(), anyhow::Error>
2004where
2005    L::TimeContainer: BatchContainer<Owned = Timestamp>,
2006    Write: AsyncFnMut(&OrdValBatch<L>, BatchDescription, &mut W) -> Result<(), anyhow::Error>,
2007    Close: AsyncFnMut(BatchDescription, &mut W) -> Result<(), anyhow::Error>,
2008{
2009    loop {
2010        {
2011            // Drop any input batches that fall below the lowest output batch.
2012            // No future output batch could need these inputs.
2013            let output_lower = output_batches
2014                .front()
2015                .map_or(&output_frontier, |((lower, _), _)| lower);
2016            while input_batches
2017                .pop_front_if(|rows| PartialOrder::less_equal(rows.upper(), output_lower))
2018                .is_some()
2019            {}
2020        }
2021
2022        {
2023            // Close and drop any output batches that fall below the lowest input batch.
2024            // No future inputs can arrive for these batches.
2025            let input_lower = input_batches
2026                .front()
2027                .map_or(&input_frontier, |rows| rows.lower());
2028            while let Some((batch_desc, mut batch_writer)) =
2029                output_batches.pop_front_if(|((_, batch_upper), _)| {
2030                    PartialOrder::less_equal(batch_upper, input_lower)
2031                })
2032            {
2033                close_batch(batch_desc, &mut batch_writer).await?;
2034            }
2035        }
2036
2037        let Some((batch_desc, batch_writer)) = output_batches.front_mut() else {
2038            // We're still waiting for descriptions of batches to write to.
2039            break;
2040        };
2041
2042        let Some(rows) = input_batches.front() else {
2043            // We're still waiting for rows to write.
2044            break;
2045        };
2046
2047        // If there were no overlap between the lowest input batch and the lowest output batch,
2048        // we'd have dropped the lower one already.
2049        // Since we still have both a lowest input batch and a lowest output batch,
2050        // there must be overlap.
2051
2052        // Write (the relevant portion of) the lowest input batch to the lowest output batch.
2053        // Drop whichever one's "upper" comes first. If they end simultaneously, drop both.
2054        write_rows(rows, batch_desc.clone(), batch_writer).await?;
2055        let output_upper = batch_desc.1.clone();
2056        let rows_upper = rows.upper();
2057        if PartialOrder::less_equal(&output_upper, rows_upper) {
2058            // Close and drop the output batch.
2059            let (batch_desc, mut batch_writer) =
2060                output_batches.pop_front().expect("already checked front");
2061            close_batch(batch_desc, &mut batch_writer).await?;
2062        }
2063        if PartialOrder::less_equal(rows_upper, &output_upper) {
2064            // Drop the input batch.
2065            input_batches.pop_front();
2066        }
2067
2068        // At least one of the two conditions above must be true,
2069        // so every loop iteration shrinks working set (of input/output batches).
2070        // Therefore, this loop must terminate.
2071    }
2072
2073    Ok(())
2074}
2075
2076#[cfg(test)]
2077mod tests {
2078    use iceberg::spec::{ListType, PrimitiveType, Type};
2079    use iceberg::writer::file_writer::location_generator::LocationGenerator;
2080    use mz_repr::SqlScalarType;
2081    use mz_storage_types::sinks::ICEBERG_UINT64_DECIMAL_PRECISION;
2082
2083    use super::*;
2084
2085    /// The URI a data file is committed under, as the manifest records it.
2086    fn manifest_uri(configured_path: Option<&str>, location: &str) -> String {
2087        let data_location = data_file_location(configured_path, location);
2088        DefaultLocationGenerator::with_data_location(data_location)
2089            .generate_location(None, "part-00000.parquet")
2090    }
2091
2092    /// Asserts the URI addresses exactly one object, i.e. it survives the path normalization
2093    /// the object store applies before writing. An empty path segment would make the manifest
2094    /// name a key that was never written.
2095    fn assert_addresses_one_object(uri: &str) {
2096        let path = uri
2097            .split_once("://")
2098            .map(|(_scheme, path)| path)
2099            .unwrap_or(uri);
2100        assert!(
2101            !path.contains("//"),
2102            "URI has an empty path segment, so it does not name the object written: {uri}"
2103        );
2104    }
2105
2106    #[mz_ore::test]
2107    fn test_data_file_location_trims_configured_path() {
2108        // The property the catalog set is honored as-is when it carries no trailing separator.
2109        assert_eq!(
2110            manifest_uri(Some("s3://bucket/tbl/data"), "s3://bucket/tbl"),
2111            "s3://bucket/tbl/data/part-00000.parquet"
2112        );
2113
2114        // A trailing separator is valid in the property, and must not reach the manifest.
2115        for configured in [
2116            "s3://bucket/tbl/data/",
2117            "s3://bucket/tbl/data//",
2118            "s3://bucket/tbl/data///",
2119        ] {
2120            let uri = manifest_uri(Some(configured), "s3://bucket/tbl");
2121            assert_addresses_one_object(&uri);
2122            assert_eq!(uri, "s3://bucket/tbl/data/part-00000.parquet");
2123        }
2124    }
2125
2126    #[mz_ore::test]
2127    fn test_data_file_location_trims_table_location() {
2128        // With no property set, the data directory hangs off the table location.
2129        assert_eq!(
2130            manifest_uri(None, "s3://bucket/tbl"),
2131            "s3://bucket/tbl/data/part-00000.parquet"
2132        );
2133
2134        // A table location ending in `/` would otherwise double the separator mid-URI, where
2135        // trimming the end of the joined string could not fix it.
2136        let uri = manifest_uri(None, "s3://bucket/tbl/");
2137        assert_addresses_one_object(&uri);
2138        assert_eq!(uri, "s3://bucket/tbl/data/part-00000.parquet");
2139    }
2140
2141    #[mz_ore::test]
2142    fn test_data_file_location_corrects_s3_tables_metadata_path() {
2143        // S3 Tables reports the metadata file as the table location; the data directory has to
2144        // hang off the warehouse root instead.
2145        assert_eq!(
2146            data_file_location(None, "s3://bucket/tbl/metadata/00001-abc.metadata.json"),
2147            "s3://bucket/tbl/data"
2148        );
2149
2150        // A path that merely contains `/metadata/` is not a metadata file and is left alone.
2151        assert_eq!(
2152            data_file_location(None, "s3://bucket/metadata/tbl"),
2153            "s3://bucket/metadata/tbl/data"
2154        );
2155    }
2156
2157    #[mz_ore::test]
2158    fn test_iceberg_type_overrides() {
2159        // UInt16 should override to Int32
2160        let result = iceberg_type_overrides(&SqlScalarType::UInt16);
2161        assert_eq!(result.unwrap().0, DataType::Int32);
2162
2163        // UInt32 should override to Int64
2164        let result = iceberg_type_overrides(&SqlScalarType::UInt32);
2165        assert_eq!(result.unwrap().0, DataType::Int64);
2166
2167        // UInt64 should override to Decimal128(20, 0)
2168        let result = iceberg_type_overrides(&SqlScalarType::UInt64);
2169        assert_eq!(
2170            result.unwrap().0,
2171            DataType::Decimal128(ICEBERG_UINT64_DECIMAL_PRECISION, 0)
2172        );
2173
2174        // MzTimestamp should override to Decimal128(20, 0)
2175        let result = iceberg_type_overrides(&SqlScalarType::MzTimestamp);
2176        assert_eq!(
2177            result.unwrap().0,
2178            DataType::Decimal128(ICEBERG_UINT64_DECIMAL_PRECISION, 0)
2179        );
2180
2181        // Interval should override to LargeUtf8
2182        let result = iceberg_type_overrides(&SqlScalarType::Interval);
2183        assert_eq!(result.unwrap().0, DataType::LargeUtf8);
2184
2185        // Uuid should override to Utf8
2186        let result = iceberg_type_overrides(&SqlScalarType::Uuid);
2187        assert_eq!(result.unwrap().0, DataType::Utf8);
2188
2189        // Other types should return None (use default)
2190        assert!(iceberg_type_overrides(&SqlScalarType::Int32).is_none());
2191        assert!(iceberg_type_overrides(&SqlScalarType::String).is_none());
2192        assert!(iceberg_type_overrides(&SqlScalarType::Bool).is_none());
2193    }
2194
2195    #[mz_ore::test]
2196    fn test_iceberg_schema_with_nested_uint64() {
2197        // Test that desc_to_schema_with_overrides handles nested UInt64
2198        // by using iceberg_type_overrides which applies recursively
2199        let desc = mz_repr::RelationDesc::builder()
2200            .with_column(
2201                "items",
2202                SqlScalarType::List {
2203                    element_type: Box::new(SqlScalarType::UInt64),
2204                    custom_id: None,
2205                }
2206                .nullable(true),
2207            )
2208            .finish();
2209
2210        let schema =
2211            mz_arrow_util::builder::desc_to_schema_with_overrides(&desc, iceberg_type_overrides)
2212                .expect("schema conversion should succeed");
2213
2214        // The inner element should be Decimal128, not UInt64
2215        if let DataType::List(field) = schema.field(0).data_type() {
2216            assert_eq!(
2217                field.data_type(),
2218                &DataType::Decimal128(ICEBERG_UINT64_DECIMAL_PRECISION, 0)
2219            );
2220        } else {
2221            panic!("Expected List type");
2222        }
2223    }
2224
2225    #[mz_ore::test]
2226    fn test_iceberg_interval_override() {
2227        // Interval should override to LargeUtf8 (string) for Iceberg
2228        let result = iceberg_type_overrides(&SqlScalarType::Interval);
2229        assert_eq!(result.unwrap().0, DataType::LargeUtf8);
2230
2231        // Test full schema conversion with interval column
2232        let desc = mz_repr::RelationDesc::builder()
2233            .with_column("id", SqlScalarType::Int32.nullable(false))
2234            .with_column("dur", SqlScalarType::Interval.nullable(true))
2235            .finish();
2236
2237        let (arrow_schema, iceberg_schema) =
2238            relation_desc_to_iceberg_schema(&desc).expect("schema conversion should succeed");
2239
2240        // Arrow schema should have LargeUtf8 for interval
2241        assert_eq!(arrow_schema.field(1).data_type(), &DataType::LargeUtf8);
2242
2243        // Iceberg schema should have String type
2244        let field = iceberg_schema
2245            .as_struct()
2246            .field_by_name("dur")
2247            .expect("field should exist");
2248        assert_eq!(*field.field_type, Type::Primitive(PrimitiveType::String));
2249    }
2250
2251    /// A uuid column must reach the Iceberg table as `string`, not as the
2252    /// `fixed[16]` that the default `FixedSizeBinary(16)` mapping would produce.
2253    #[mz_ore::test]
2254    fn test_iceberg_uuid_override() {
2255        let result = iceberg_type_overrides(&SqlScalarType::Uuid);
2256        assert_eq!(result.unwrap().0, DataType::Utf8);
2257
2258        let desc = mz_repr::RelationDesc::builder()
2259            .with_column("id", SqlScalarType::Int32.nullable(false))
2260            .with_column("u", SqlScalarType::Uuid.nullable(true))
2261            .finish();
2262
2263        let (arrow_schema, iceberg_schema) =
2264            relation_desc_to_iceberg_schema(&desc).expect("schema conversion should succeed");
2265
2266        assert_eq!(arrow_schema.field(1).data_type(), &DataType::Utf8);
2267
2268        let field = iceberg_schema
2269            .as_struct()
2270            .field_by_name("u")
2271            .expect("field should exist");
2272        assert_eq!(*field.field_type, Type::Primitive(PrimitiveType::String));
2273        assert_ne!(
2274            *field.field_type,
2275            Type::Primitive(PrimitiveType::Fixed(16)),
2276            "uuid must not fall back to the default FixedSizeBinary(16) mapping"
2277        );
2278    }
2279
2280    /// Rebuilds `schema` with the named field's type replaced, so a fixture can
2281    /// differ from what the sink derives in exactly one type.
2282    fn with_field_type(schema: &Schema, name: &str, ty: Type) -> Schema {
2283        let fields = schema.as_struct().fields().iter().map(|f| {
2284            let mut field = (**f).clone();
2285            if field.name == name {
2286                field.field_type = Box::new(ty.clone());
2287            }
2288            Arc::new(field)
2289        });
2290        Schema::builder()
2291            .with_fields(fields)
2292            .build()
2293            .expect("valid schema")
2294    }
2295
2296    /// A table created before uuid columns became strings holds `fixed[16]` and
2297    /// must stay writable, while any other difference is still a mismatch.
2298    #[mz_ore::test]
2299    fn test_is_compatible_accepts_legacy_uuid() {
2300        let desc = mz_repr::RelationDesc::builder()
2301            .with_column("id", SqlScalarType::Int32.nullable(false))
2302            .with_column("u", SqlScalarType::Uuid.nullable(true))
2303            .finish();
2304        let (_, expected) =
2305            relation_desc_to_iceberg_schema(&desc).expect("schema conversion should succeed");
2306
2307        assert!(is_compatible(&expected, &expected));
2308
2309        // The halt this tolerance exists to prevent.
2310        let legacy = with_field_type(&expected, "u", Type::Primitive(PrimitiveType::Fixed(16)));
2311        assert!(is_compatible(&legacy, &expected));
2312
2313        // Only that direction: nothing produces a string column where the sink
2314        // derives fixed binary.
2315        assert!(!is_compatible(&expected, &legacy));
2316
2317        // An unrelated type difference is still a mismatch.
2318        let wrong_type = with_field_type(&expected, "u", Type::Primitive(PrimitiveType::Long));
2319        assert!(!is_compatible(&wrong_type, &expected));
2320
2321        // So is a missing column.
2322        let truncated = Schema::builder()
2323            .with_fields(expected.as_struct().fields().iter().take(1).cloned())
2324            .build()
2325            .expect("valid schema");
2326        assert!(!is_compatible(&truncated, &expected));
2327    }
2328
2329    /// The tolerance reaches a uuid nested inside a list, which the type
2330    /// overrides remap just like a top-level column.
2331    #[mz_ore::test]
2332    fn test_is_compatible_accepts_legacy_uuid_in_list() {
2333        let desc = mz_repr::RelationDesc::builder()
2334            .with_column("id", SqlScalarType::Int32.nullable(false))
2335            .with_column(
2336                "us",
2337                SqlScalarType::List {
2338                    element_type: Box::new(SqlScalarType::Uuid),
2339                    custom_id: None,
2340                }
2341                .nullable(true),
2342            )
2343            .finish();
2344        let (_, expected) =
2345            relation_desc_to_iceberg_schema(&desc).expect("schema conversion should succeed");
2346
2347        let Type::List(list) = &*expected
2348            .as_struct()
2349            .field_by_name("us")
2350            .expect("field should exist")
2351            .field_type
2352        else {
2353            panic!("expected a list type");
2354        };
2355        let mut legacy_element = (*list.element_field).clone();
2356        legacy_element.field_type = Box::new(Type::Primitive(PrimitiveType::Fixed(16)));
2357        let legacy = with_field_type(
2358            &expected,
2359            "us",
2360            Type::List(ListType {
2361                element_field: Arc::new(legacy_element),
2362            }),
2363        );
2364
2365        assert!(is_compatible(&legacy, &expected));
2366    }
2367
2368    #[mz_ore::test]
2369    fn test_iceberg_range_schema() {
2370        // Test full schema conversion with range column
2371        let desc = mz_repr::RelationDesc::builder()
2372            .with_column("id", SqlScalarType::Int32.nullable(false))
2373            .with_column(
2374                "r",
2375                SqlScalarType::Range {
2376                    element_type: Box::new(SqlScalarType::Int32),
2377                }
2378                .nullable(true),
2379            )
2380            .finish();
2381
2382        let (_arrow_schema, iceberg_schema) =
2383            relation_desc_to_iceberg_schema(&desc).expect("schema conversion should succeed");
2384
2385        // Iceberg schema should have a struct type for the range
2386        let field = iceberg_schema
2387            .as_struct()
2388            .field_by_name("r")
2389            .expect("field should exist");
2390        assert!(
2391            matches!(&*field.field_type, Type::Struct(_)),
2392            "range should be struct, got: {:?}",
2393            field.field_type
2394        );
2395    }
2396
2397    #[mz_ore::test]
2398    fn equality_ids_follow_iceberg_field_ids() {
2399        let map_entries = Field::new(
2400            "entries",
2401            DataType::Struct(
2402                vec![
2403                    Field::new("key", DataType::Utf8, false),
2404                    Field::new("value", DataType::Utf8, true),
2405                ]
2406                .into(),
2407            ),
2408            false,
2409        );
2410        let materialize_arrow_schema = ArrowSchema::new(vec![
2411            Field::new("attrs", DataType::Map(Arc::new(map_entries), false), true),
2412            Field::new("key_col", DataType::Int32, false),
2413        ]);
2414        let materialize_arrow_schema = add_field_ids_to_arrow_schema(materialize_arrow_schema);
2415        let iceberg_schema = arrow_schema_to_schema(&materialize_arrow_schema)
2416            .expect("schema conversion should succeed");
2417
2418        let equality_ids =
2419            equality_ids_for_indices(&iceberg_schema, &materialize_arrow_schema, &[1])
2420                .expect("field lookup should succeed");
2421
2422        let expected_id = iceberg_schema
2423            .as_struct()
2424            .field_by_name("key_col")
2425            .expect("top-level field should exist")
2426            .id;
2427        assert_eq!(equality_ids, vec![expected_id]);
2428        assert_ne!(expected_id, 2);
2429    }
2430
2431    /// Regression test: iceberg-rust names map fields `key_value`/`key`/`value`
2432    /// while Materialize uses `entries`/`keys`/`values`. The schema merge must
2433    /// still copy the value field's extension metadata across so ArrowBuilder
2434    /// can build the inner builder.
2435    #[mz_ore::test]
2436    #[allow(clippy::disallowed_types)]
2437    fn merge_map_entries_preserves_value_extension_metadata() {
2438        use std::collections::HashMap;
2439
2440        let mz_value_metadata = HashMap::from([(
2441            ARROW_EXTENSION_NAME_KEY.to_string(),
2442            "materialize.v1.string".to_string(),
2443        )]);
2444        let mz_entries = Field::new(
2445            "entries",
2446            DataType::Struct(
2447                vec![
2448                    Field::new("keys", DataType::Utf8, false),
2449                    Field::new("values", DataType::Utf8, true).with_metadata(mz_value_metadata),
2450                ]
2451                .into(),
2452            ),
2453            false,
2454        );
2455        let mz_map = Field::new("m", DataType::Map(Arc::new(mz_entries), false), true)
2456            .with_metadata(HashMap::from([(
2457                ARROW_EXTENSION_NAME_KEY.to_string(),
2458                "materialize.v1.map".to_string(),
2459            )]));
2460
2461        let iceberg_entries = Field::new(
2462            "key_value",
2463            DataType::Struct(
2464                vec![
2465                    Field::new("key", DataType::Utf8, false),
2466                    Field::new("value", DataType::Utf8, true),
2467                ]
2468                .into(),
2469            ),
2470            false,
2471        );
2472        let iceberg_map = Field::new("m", DataType::Map(Arc::new(iceberg_entries), false), true);
2473
2474        let merged = merge_field_metadata_recursive(&iceberg_map, Some(&mz_map))
2475            .expect("merge should succeed");
2476
2477        let entries = match merged.data_type() {
2478            DataType::Map(entries, _) => entries.as_ref(),
2479            other => panic!("expected Map, got {other:?}"),
2480        };
2481        let entry_fields = match entries.data_type() {
2482            DataType::Struct(fields) => fields,
2483            other => panic!("expected Struct, got {other:?}"),
2484        };
2485        // Iceberg naming must be preserved on the merged schema...
2486        assert_eq!(entry_fields[0].name(), "key");
2487        assert_eq!(entry_fields[1].name(), "value");
2488        // ...and the materialize extension must have been copied positionally
2489        // to the value field even though its name didn't match `values`.
2490        assert_eq!(
2491            entry_fields[1].metadata().get(ARROW_EXTENSION_NAME_KEY),
2492            Some(&"materialize.v1.string".to_string()),
2493        );
2494    }
2495
2496    mod with_ready_batches {
2497        use differential_dataflow::trace::Batch;
2498        use differential_dataflow::trace::implementations::Vector;
2499
2500        use super::*;
2501
2502        type TestBatch = OrdValBatch<Vector<((u64, u64), Timestamp, Diff)>>;
2503
2504        /// A frontier at `t`, or the empty (end-of-time) frontier for `None`.
2505        fn frontier(t: Option<u64>) -> Antichain<Timestamp> {
2506            t.map_or_else(Antichain::new, |t| Antichain::from_elem(Timestamp::new(t)))
2507        }
2508
2509        /// `[lower, upper)` bounds, with `None` for the empty upper.
2510        fn span(lower: u64, upper: Option<u64>) -> BatchDescription {
2511            (frontier(Some(lower)), frontier(upper))
2512        }
2513
2514        /// An input batch with the given bounds. The pairing logic under test
2515        /// only looks at bounds, so the batch holds no data.
2516        fn input(lower: u64, upper: Option<u64>) -> ArcBatch<TestBatch> {
2517            let (lower, upper) = span(lower, upper);
2518            ArcBatch(Arc::new(TestBatch::empty(lower, upper)))
2519        }
2520
2521        #[derive(Debug, PartialEq)]
2522        enum Call {
2523            /// (input batch bounds, output batch description)
2524            Write(BatchDescription, BatchDescription),
2525            Close(BatchDescription),
2526        }
2527
2528        /// Run `with_ready_batches` with recording callbacks and return the
2529        /// sequence of calls it made.
2530        async fn run(
2531            input_frontier: Antichain<Timestamp>,
2532            input_batches: &mut VecDeque<ArcBatch<TestBatch>>,
2533            output_frontier: Antichain<Timestamp>,
2534            output_batches: &mut VecDeque<(BatchDescription, ())>,
2535        ) -> Vec<Call> {
2536            let calls = RefCell::new(vec![]);
2537            with_ready_batches(
2538                input_frontier,
2539                input_batches,
2540                output_frontier,
2541                output_batches,
2542                async |rows: &TestBatch, desc, _writer: &mut ()| {
2543                    let bounds = (rows.lower().clone(), rows.upper().clone());
2544                    calls.borrow_mut().push(Call::Write(bounds, desc));
2545                    Ok(())
2546                },
2547                async |desc, _writer: &mut ()| {
2548                    calls.borrow_mut().push(Call::Close(desc));
2549                    Ok(())
2550                },
2551            )
2552            .await
2553            .expect("test callbacks never fail");
2554            calls.into_inner()
2555        }
2556
2557        #[mz_ore::test(tokio::test)]
2558        async fn input_batch_spanning_multiple_output_batches() {
2559            let mut inputs = VecDeque::from([input(0, Some(30))]);
2560            let mut outputs = VecDeque::from([
2561                (span(0, Some(10)), ()),
2562                (span(10, Some(20)), ()),
2563                (span(20, Some(30)), ()),
2564            ]);
2565
2566            let calls = run(
2567                frontier(Some(30)),
2568                &mut inputs,
2569                frontier(Some(30)),
2570                &mut outputs,
2571            )
2572            .await;
2573
2574            // The input batch is written once per overlapping output batch,
2575            // each of which closes as soon as the input covers its upper.
2576            assert_eq!(
2577                calls,
2578                vec![
2579                    Call::Write(span(0, Some(30)), span(0, Some(10))),
2580                    Call::Close(span(0, Some(10))),
2581                    Call::Write(span(0, Some(30)), span(10, Some(20))),
2582                    Call::Close(span(10, Some(20))),
2583                    Call::Write(span(0, Some(30)), span(20, Some(30))),
2584                    Call::Close(span(20, Some(30))),
2585                ]
2586            );
2587            assert!(inputs.is_empty());
2588            assert!(outputs.is_empty());
2589        }
2590
2591        #[mz_ore::test(tokio::test)]
2592        async fn output_batch_spanning_multiple_input_batches() {
2593            let mut inputs =
2594                VecDeque::from([input(0, Some(10)), input(10, Some(20)), input(20, Some(30))]);
2595            let mut outputs = VecDeque::from([(span(0, Some(30)), ())]);
2596
2597            let calls = run(
2598                frontier(Some(30)),
2599                &mut inputs,
2600                frontier(Some(30)),
2601                &mut outputs,
2602            )
2603            .await;
2604
2605            assert_eq!(
2606                calls,
2607                vec![
2608                    Call::Write(span(0, Some(10)), span(0, Some(30))),
2609                    Call::Write(span(10, Some(20)), span(0, Some(30))),
2610                    Call::Write(span(20, Some(30)), span(0, Some(30))),
2611                    Call::Close(span(0, Some(30))),
2612                ]
2613            );
2614            assert!(inputs.is_empty());
2615            assert!(outputs.is_empty());
2616        }
2617
2618        #[mz_ore::test(tokio::test)]
2619        async fn input_batch_retained_for_future_output_batches() {
2620            let mut inputs = VecDeque::from([input(0, Some(30))]);
2621            let mut outputs = VecDeque::from([(span(0, Some(10)), ())]);
2622
2623            let calls = run(
2624                frontier(Some(30)),
2625                &mut inputs,
2626                frontier(Some(10)),
2627                &mut outputs,
2628            )
2629            .await;
2630
2631            // The input batch extends past the only known output batch, so it
2632            // must stay queued for descriptions that haven't arrived yet.
2633            assert_eq!(
2634                calls,
2635                vec![
2636                    Call::Write(span(0, Some(30)), span(0, Some(10))),
2637                    Call::Close(span(0, Some(10))),
2638                ]
2639            );
2640            assert_eq!(inputs.len(), 1);
2641            assert!(outputs.is_empty());
2642        }
2643
2644        #[mz_ore::test(tokio::test)]
2645        async fn already_committed_input_batches_dropped_unwritten() {
2646            let mut inputs = VecDeque::from([input(0, Some(10)), input(10, Some(20))]);
2647            let mut outputs = VecDeque::from([(span(20, Some(30)), ())]);
2648
2649            let calls = run(
2650                frontier(Some(20)),
2651                &mut inputs,
2652                frontier(Some(30)),
2653                &mut outputs,
2654            )
2655            .await;
2656
2657            // Both input batches fall below the lowest output batch, so their
2658            // contents are already committed and they are dropped unwritten.
2659            // The output batch still waits for its own input.
2660            assert_eq!(calls, vec![]);
2661            assert!(inputs.is_empty());
2662            assert_eq!(outputs.len(), 1);
2663        }
2664
2665        #[mz_ore::test(tokio::test)]
2666        async fn output_batch_closes_empty_once_input_frontier_passes() {
2667            let mut outputs = VecDeque::from([(span(0, Some(10)), ())]);
2668
2669            // While the input frontier is short of the batch's upper, nothing
2670            // may close: rows for it could still arrive.
2671            let calls = run(
2672                frontier(Some(5)),
2673                &mut VecDeque::new(),
2674                frontier(Some(10)),
2675                &mut outputs,
2676            )
2677            .await;
2678            assert_eq!(calls, vec![]);
2679            assert_eq!(outputs.len(), 1);
2680
2681            // Once the input frontier reaches the upper, the batch closes
2682            // empty (an empty commit).
2683            let calls = run(
2684                frontier(Some(10)),
2685                &mut VecDeque::new(),
2686                frontier(Some(10)),
2687                &mut outputs,
2688            )
2689            .await;
2690            assert_eq!(calls, vec![Call::Close(span(0, Some(10)))]);
2691            assert!(outputs.is_empty());
2692        }
2693
2694        #[mz_ore::test(tokio::test)]
2695        async fn final_output_batch_with_empty_upper() {
2696            let mut inputs = VecDeque::from([input(20, Some(30))]);
2697            let mut outputs = VecDeque::from([(span(20, None), ())]);
2698
2699            // The sealing batch covers everything from 20 to the end of time.
2700            // It consumes all remaining input but only closes once the input
2701            // frontier is empty, i.e. the input is finished.
2702            let calls = run(
2703                frontier(Some(30)),
2704                &mut inputs,
2705                frontier(None),
2706                &mut outputs,
2707            )
2708            .await;
2709            assert_eq!(calls, vec![Call::Write(span(20, Some(30)), span(20, None))]);
2710            assert!(inputs.is_empty());
2711            assert_eq!(outputs.len(), 1);
2712
2713            let calls = run(frontier(None), &mut inputs, frontier(None), &mut outputs).await;
2714            assert_eq!(calls, vec![Call::Close(span(20, None))]);
2715            assert!(outputs.is_empty());
2716        }
2717    }
2718}
2719
2720/// Commit completed batches to Iceberg as snapshots.
2721/// Batches are committed in timestamp order to ensure strong consistency guarantees downstream.
2722/// Each snapshot includes the Materialize frontier in its metadata for resume support.
2723fn commit_to_iceberg<'scope>(
2724    name: String,
2725    sink_id: GlobalId,
2726    sink_version: u64,
2727    batch_input: StreamVec<'scope, Timestamp, BoundedDataFile>,
2728    batch_desc_input: StreamVec<'scope, Timestamp, (Antichain<Timestamp>, Antichain<Timestamp>)>,
2729    table_ready_stream: StreamVec<'scope, Timestamp, Infallible>,
2730    write_frontier: Rc<RefCell<Antichain<Timestamp>>>,
2731    connection: IcebergSinkConnection,
2732    storage_configuration: StorageConfiguration,
2733    write_handle: impl Future<
2734        Output = anyhow::Result<WriteHandle<SourceData, (), Timestamp, StorageDiff>>,
2735    > + 'static,
2736    metrics: Arc<IcebergSinkMetrics>,
2737    statistics: SinkStatistics,
2738) -> (
2739    StreamVec<'scope, Timestamp, HealthStatusMessage>,
2740    PressOnDropButton,
2741) {
2742    let scope = batch_input.scope();
2743    let mut builder = OperatorBuilder::new(name, scope.clone());
2744
2745    let hashed_id = sink_id.hashed();
2746    let is_active_worker = usize::cast_from(hashed_id) % scope.peers() == scope.index();
2747    let name_for_logging = format!("{sink_id}-commit-to-iceberg");
2748
2749    let mut input = builder.new_disconnected_input(batch_input, Exchange::new(move |_| hashed_id));
2750    let mut batch_desc_input =
2751        builder.new_disconnected_input(batch_desc_input, Exchange::new(move |_| hashed_id));
2752    let mut table_ready_input = builder.new_disconnected_input(table_ready_stream, Pipeline);
2753
2754    let (button, errors) = builder.build_fallible(move |_caps| {
2755        Box::pin(async move {
2756            if !is_active_worker {
2757                write_frontier.borrow_mut().clear();
2758                return Ok(());
2759            }
2760
2761            let namespace_ident = NamespaceIdent::new(connection.namespace.clone());
2762            let table_ident = TableIdent::new(namespace_ident, connection.table.clone());
2763            let catalog = connection
2764                .catalog_connection
2765                .connect(&storage_configuration, InTask::Yes, Some(&table_ident))
2766                .await
2767                .with_context(|| {
2768                    format!(
2769                        "Failed to connect to Iceberg catalog '{}' for table '{}.{}'",
2770                        connection.catalog_connection.uri, connection.namespace, connection.table
2771                    )
2772                })?;
2773
2774            let mut write_handle = write_handle.await?;
2775
2776            while let Some(_) = table_ready_input.next().await {
2777                // Wait for table to be ready
2778            }
2779            let mut table = catalog.load_table(&table_ident).await.with_context(|| {
2780                format!(
2781                    "Failed to load Iceberg table '{}.{}' in commit_to_iceberg operator",
2782                    connection.namespace, connection.table
2783                )
2784            })?;
2785
2786            #[allow(clippy::disallowed_types)]
2787            let mut batch_descriptions: std::collections::HashMap<
2788                (Antichain<Timestamp>, Antichain<Timestamp>),
2789                BoundedDataFileSet,
2790            > = std::collections::HashMap::new();
2791
2792            let mut batch_description_frontier = Antichain::from_elem(Timestamp::minimum());
2793            let mut input_frontier = Antichain::from_elem(Timestamp::minimum());
2794
2795            while !(batch_description_frontier.is_empty() && input_frontier.is_empty()) {
2796                tokio::select! {
2797                    _ = batch_desc_input.ready() => {},
2798                    _ = input.ready() => {}
2799                }
2800
2801                while let Some(event) = batch_desc_input.next_sync() {
2802                    match event {
2803                        Event::Data(_cap, data) => {
2804                            for batch_desc in data {
2805                                let prev = batch_descriptions
2806                                    .insert(batch_desc, BoundedDataFileSet { data_files: vec![] });
2807                                if let Some(prev) = prev {
2808                                    anyhow::bail!(
2809                                        "Duplicate batch description received \
2810                                         in commit operator: {:?}",
2811                                        prev
2812                                    );
2813                                }
2814                            }
2815                        }
2816                        Event::Progress(frontier) => {
2817                            batch_description_frontier = frontier;
2818                        }
2819                    }
2820                }
2821
2822                let ready_events = std::iter::from_fn(|| input.next_sync()).collect_vec();
2823                for event in ready_events {
2824                    match event {
2825                        Event::Data(_cap, data) => {
2826                            for bounded_data_file in data {
2827                                let entry = batch_descriptions
2828                                    .entry(bounded_data_file.batch_desc().clone())
2829                                    .or_default();
2830                                entry.data_files.push(bounded_data_file);
2831                            }
2832                        }
2833                        Event::Progress(frontier) => {
2834                            input_frontier = frontier;
2835                        }
2836                    }
2837                }
2838
2839                // Collect batches whose data files have all arrived.
2840                // The writer emits all data files for a batch at a capability <= the batch's
2841                // lower bound, then downgrades its capability to the batch's upper bound.
2842                // So once the input frontier advances past lower, we know the writer has
2843                // finished emitting files for this batch and dropped its capability.
2844                let mut done_batches: Vec<_> = batch_descriptions
2845                    .keys()
2846                    .filter(|(lower, _upper)| PartialOrder::less_than(lower, &input_frontier))
2847                    .cloned()
2848                    .collect();
2849
2850                // Commit batches in timestamp order to maintain consistency
2851                done_batches.sort_by(|a, b| {
2852                    if PartialOrder::less_than(a, b) {
2853                        Ordering::Less
2854                    } else if PartialOrder::less_than(b, a) {
2855                        Ordering::Greater
2856                    } else {
2857                        Ordering::Equal
2858                    }
2859                });
2860
2861                for batch in done_batches {
2862                    let file_set = batch_descriptions.remove(&batch).unwrap();
2863
2864                    let mut data_files = vec![];
2865                    let mut delete_files = vec![];
2866                    // Track totals for committed statistics
2867                    let mut total_messages: u64 = 0;
2868                    let mut total_bytes: u64 = 0;
2869                    for file in file_set.data_files {
2870                        total_messages += file.data_file().record_count();
2871                        total_bytes += file.data_file().file_size_in_bytes();
2872                        match file.data_file().content_type() {
2873                            iceberg::spec::DataContentType::Data => {
2874                                data_files.push(file.into_data_file());
2875                            }
2876                            iceberg::spec::DataContentType::PositionDeletes
2877                            | iceberg::spec::DataContentType::EqualityDeletes => {
2878                                delete_files.push(file.into_data_file());
2879                            }
2880                        }
2881                    }
2882
2883                    debug!(
2884                        ?sink_id,
2885                        %name_for_logging,
2886                        lower = %batch.0.pretty(),
2887                        upper = %batch.1.pretty(),
2888                        data_files = data_files.len(),
2889                        delete_files = delete_files.len(),
2890                        total_messages,
2891                        total_bytes,
2892                        "iceberg commit applying batch"
2893                    );
2894
2895                    let instant = Instant::now();
2896
2897                    let frontier = batch.1.clone();
2898                    let frontier_json = serde_json::to_string(&frontier.elements())
2899                        .context("Failed to serialize frontier to JSON")?;
2900                    let snapshot_properties = vec![
2901                        ("mz-sink-id".to_string(), sink_id.to_string()),
2902                        ("mz-frontier".to_string(), frontier_json),
2903                        ("mz-sink-version".to_string(), sink_version.to_string()),
2904                    ];
2905
2906                    let (table_state, commit_result) = Retry::default()
2907                        .max_tries(5)
2908                        .retry_async_with_state(table, |_, table| {
2909                            let snapshot_properties = snapshot_properties.clone();
2910                            let data_files = data_files.clone();
2911                            let delete_files = delete_files.clone();
2912                            let metrics = Arc::clone(&metrics);
2913                            let catalog = Arc::clone(&catalog);
2914                            let conn_namespace = connection.namespace.clone();
2915                            let conn_table = connection.table.clone();
2916                            let batch_lower = batch.0.clone();
2917                            let batch_upper = batch.1.clone();
2918                            async move {
2919                                try_commit_batch(
2920                                    table,
2921                                    snapshot_properties,
2922                                    data_files,
2923                                    delete_files,
2924                                    catalog.as_ref(),
2925                                    &conn_namespace,
2926                                    &conn_table,
2927                                    sink_id,
2928                                    sink_version,
2929                                    &batch_lower,
2930                                    &batch_upper,
2931                                    &metrics,
2932                                )
2933                                .await
2934                            }
2935                        })
2936                        .await;
2937                    let commit_result = commit_result.with_context(|| {
2938                        format!(
2939                            "failed to commit batch to Iceberg table '{}.{}'",
2940                            connection.namespace, connection.table
2941                        )
2942                    });
2943                    table = table_state;
2944                    let duration = instant.elapsed();
2945                    metrics
2946                        .commit_duration_seconds
2947                        .observe(duration.as_secs_f64());
2948                    commit_result?;
2949
2950                    debug!(
2951                        ?sink_id,
2952                        %name_for_logging,
2953                        lower = %batch.0.pretty(),
2954                        upper = %batch.1.pretty(),
2955                        total_messages,
2956                        total_bytes,
2957                        ?duration,
2958                        "iceberg commit applied batch"
2959                    );
2960
2961                    metrics.snapshots_committed.inc();
2962                    statistics.inc_messages_committed_by(total_messages);
2963                    statistics.inc_bytes_committed_by(total_bytes);
2964
2965                    let mut expect_upper = write_handle.shared_upper();
2966                    loop {
2967                        if PartialOrder::less_equal(&frontier, &expect_upper) {
2968                            // The frontier has already been advanced as far as necessary.
2969                            break;
2970                        }
2971
2972                        const EMPTY: &[((SourceData, ()), Timestamp, StorageDiff)] = &[];
2973                        match write_handle
2974                            .compare_and_append(EMPTY, expect_upper, frontier.clone())
2975                            .await
2976                            .expect("valid usage")
2977                        {
2978                            Ok(()) => break,
2979                            Err(mismatch) => {
2980                                expect_upper = mismatch.current;
2981                            }
2982                        }
2983                    }
2984                    write_frontier.borrow_mut().clone_from(&frontier);
2985                }
2986            }
2987
2988            Ok(())
2989        })
2990    });
2991
2992    let statuses = errors.map(|error| HealthStatusMessage {
2993        id: None,
2994        update: HealthStatusUpdate::halting(format!("{}", error.display_with_causes()), None),
2995        namespace: StatusNamespace::Iceberg,
2996    });
2997
2998    (statuses, button.press_on_drop())
2999}
3000
3001impl<'scope> SinkRender<'scope> for IcebergSinkConnection {
3002    fn get_key_indices(&self) -> Option<&[usize]> {
3003        self.key_desc_and_indices
3004            .as_ref()
3005            .map(|(_, indices)| indices.as_slice())
3006    }
3007
3008    fn get_relation_key_indices(&self) -> Option<&[usize]> {
3009        self.relation_key_indices.as_deref()
3010    }
3011
3012    fn render_sink(
3013        &self,
3014        storage_state: &mut StorageState,
3015        sink: &StorageSinkDesc<CollectionMetadata, Timestamp>,
3016        sink_id: GlobalId,
3017        batches: SinkBatchStream<'scope>,
3018        key_is_synthetic: bool,
3019        _err_collection: VecCollection<'scope, Timestamp, DataflowError, Diff>,
3020    ) -> (
3021        StreamVec<'scope, Timestamp, HealthStatusMessage>,
3022        Vec<PressOnDropButton>,
3023    ) {
3024        let scope = batches.scope();
3025
3026        let write_handle = {
3027            let persist = Arc::clone(&storage_state.persist_clients);
3028            let shard_meta = sink.to_storage_metadata.clone();
3029            async move {
3030                let client = persist.open(shard_meta.persist_location).await?;
3031                let handle = client
3032                    .open_writer(
3033                        shard_meta.data_shard,
3034                        Arc::new(shard_meta.relation_desc),
3035                        Arc::new(UnitSchema),
3036                        Diagnostics::from_purpose("sink handle"),
3037                    )
3038                    .await?;
3039                Ok(handle)
3040            }
3041        };
3042
3043        let write_frontier = Rc::new(RefCell::new(Antichain::from_elem(Timestamp::minimum())));
3044        storage_state
3045            .sink_write_frontiers
3046            .insert(sink_id, Rc::clone(&write_frontier));
3047
3048        let (arrow_schema_with_ids, iceberg_schema) =
3049            match (|| -> Result<(ArrowSchema, Arc<Schema>), anyhow::Error> {
3050                let (arrow_schema_with_ids, iceberg_schema) =
3051                    relation_desc_to_iceberg_schema(&sink.from_desc)?;
3052
3053                Ok(if sink.envelope == SinkEnvelope::Append {
3054                    // For append mode, extend the Arrow and Iceberg schemas with the user-visible
3055                    // `_mz_diff` and `_mz_timestamp` columns. The minter uses `iceberg_schema` to create
3056                    // the Iceberg table, and `write_data_files` uses `arrow_schema_with_ids` when
3057                    // merging metadata. Both must include these columns before any operator starts.
3058                    let extended_arrow = build_schema_with_append_columns(&arrow_schema_with_ids);
3059                    let extended_iceberg = Arc::new(
3060                        arrow_schema_to_schema(&extended_arrow)
3061                            .context("Failed to build Iceberg schema with append columns")?,
3062                    );
3063                    (extended_arrow, extended_iceberg)
3064                } else {
3065                    (arrow_schema_with_ids, iceberg_schema)
3066                })
3067            })() {
3068                Ok(schemas) => schemas,
3069                Err(err) => {
3070                    let error_stream = std::iter::once(HealthStatusMessage {
3071                        id: None,
3072                        update: HealthStatusUpdate::halting(
3073                            format!("{}", err.display_with_causes()),
3074                            None,
3075                        ),
3076                        namespace: StatusNamespace::Iceberg,
3077                    })
3078                    .to_stream(scope);
3079                    return (error_stream, vec![]);
3080                }
3081            };
3082
3083        let metrics = Arc::new(
3084            storage_state
3085                .metrics
3086                .get_iceberg_sink_metrics(sink_id, scope.index()),
3087        );
3088
3089        let statistics = storage_state
3090            .aggregated_statistics
3091            .get_sink(&sink_id)
3092            .expect("statistics initialized")
3093            .clone();
3094
3095        let connection_for_minter = self.clone();
3096        let (batch_descriptions, table_ready, mint_status, mint_button) = mint_batch_descriptions(
3097            format!("{sink_id}-iceberg-mint"),
3098            sink_id,
3099            batches.clone(),
3100            sink,
3101            connection_for_minter,
3102            storage_state.storage_configuration.clone(),
3103            Arc::clone(&iceberg_schema),
3104        );
3105
3106        let connection_for_writer = self.clone();
3107        let (datafiles, write_status, write_button) = match sink.envelope {
3108            SinkEnvelope::Upsert => write_data_files::<UpsertEnvelopeHandler>(
3109                format!("{sink_id}-write-data-files"),
3110                batches,
3111                batch_descriptions.clone(),
3112                table_ready.clone(),
3113                sink_id,
3114                sink.from,
3115                key_is_synthetic,
3116                sink.as_of.clone(),
3117                connection_for_writer,
3118                storage_state.storage_configuration.clone(),
3119                Arc::new(arrow_schema_with_ids.clone()),
3120                Arc::clone(&metrics),
3121                statistics.clone(),
3122            ),
3123            SinkEnvelope::Append => write_data_files::<AppendEnvelopeHandler>(
3124                format!("{sink_id}-write-data-files"),
3125                batches,
3126                batch_descriptions.clone(),
3127                table_ready.clone(),
3128                sink_id,
3129                sink.from,
3130                key_is_synthetic,
3131                sink.as_of.clone(),
3132                connection_for_writer,
3133                storage_state.storage_configuration.clone(),
3134                Arc::new(arrow_schema_with_ids.clone()),
3135                Arc::clone(&metrics),
3136                statistics.clone(),
3137            ),
3138            SinkEnvelope::Debezium => {
3139                unreachable!("Iceberg sink only supports Upsert and Append envelopes")
3140            }
3141        };
3142
3143        let connection_for_committer = self.clone();
3144        let (commit_status, commit_button) = commit_to_iceberg(
3145            format!("{sink_id}-commit-to-iceberg"),
3146            sink_id,
3147            sink.version,
3148            datafiles,
3149            batch_descriptions,
3150            table_ready,
3151            Rc::clone(&write_frontier),
3152            connection_for_committer,
3153            storage_state.storage_configuration.clone(),
3154            write_handle,
3155            Arc::clone(&metrics),
3156            statistics,
3157        );
3158
3159        let running_status = Some(HealthStatusMessage {
3160            id: None,
3161            update: HealthStatusUpdate::running(),
3162            namespace: StatusNamespace::Iceberg,
3163        })
3164        .to_stream(scope);
3165
3166        let statuses =
3167            scope.concatenate([running_status, mint_status, write_status, commit_status]);
3168
3169        (statuses, vec![mint_button, write_button, commit_button])
3170    }
3171}