1use 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
166const DEFAULT_ARRAY_BUILDER_ITEM_CAPACITY: usize = 1024;
169const DEFAULT_ARRAY_BUILDER_DATA_CAPACITY: usize = 1024;
173
174const PARQUET_FILE_PREFIX: &str = "mz_data";
176const INITIAL_DESCRIPTIONS_TO_MINT: u64 = 3;
179
180struct WriterContext {
183 arrow_schema: Arc<ArrowSchema>,
185 current_schema: Arc<Schema>,
187 file_io: iceberg::io::FileIO,
189 location_generator: DefaultLocationGenerator,
191 file_name_generator: DefaultFileNameGenerator,
193 writer_properties: WriterProperties,
194}
195
196trait EnvelopeHandler: Send {
198 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 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 equality_ids: Vec<i32>,
221 pos_schema: Arc<Schema>,
223 eq_schema: Arc<Schema>,
225 eq_config: EqualityDeleteWriterConfig,
227 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 builder.with_max_seen_rows(0)
332 } else {
333 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 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 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 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 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
483fn 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
497fn 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
513fn 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
539fn merge_materialize_metadata_into_iceberg_schema(
544 materialize_arrow_schema: &ArrowSchema,
545 iceberg_schema: &Schema,
546) -> anyhow::Result<ArrowSchema> {
547 let iceberg_arrow_schema = schema_to_arrow_schema(iceberg_schema)
549 .context("Failed to convert Iceberg schema to Arrow schema")?;
550
551 let fields: Vec<Field> = iceberg_arrow_schema
553 .fields()
554 .iter()
555 .map(|iceberg_field| {
556 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
573fn merge_field_metadata_recursive(
575 iceberg_field: &Field,
576 mz_field: Option<&Field>,
577) -> anyhow::Result<Field> {
578 let mut metadata = iceberg_field.metadata().clone();
580
581 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 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 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
678fn 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
781enum CommitError {
783 Local(iceberg::Error),
786 Request(iceberg::Error),
789}
790
791async 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 let mut action_commit = Arc::new(action)
824 .commit(table)
825 .await
826 .map_err(CommitError::Local)?;
827
828 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
845async 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 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 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 if last_id == sink_id && last_version == sink_version && last_frontier == *batch_upper {
886 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 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 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 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 metrics.commit_conflicts.inc();
954 (table, RetryResult::RetryableErr(anyhow!(e)))
955 }
956 ErrorKind::Unexpected => {
957 metrics.commit_failures.inc();
960 (table, RetryResult::RetryableErr(anyhow!(e)))
961 }
962 _ => {
963 metrics.commit_failures.inc();
965 (table, RetryResult::FatalErr(anyhow!(e)))
966 }
967 },
968 }
969}
970
971fn 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 (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
1018async 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 match catalog.load_table(&table_ident).await {
1030 Ok(table) => {
1031 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 let table_creation = TableCreation::builder()
1055 .name(table_name.clone())
1056 .schema(schema.clone())
1057 .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 Err(err).context("Failed to load Iceberg table")
1074 }
1075 }
1076 }
1077}
1078
1079fn 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 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
1127fn 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
1151fn 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
1183fn 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#[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
1212fn 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 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 let overcompacted =
1308 *resume_upper != [Timestamp::minimum()] &&
1310 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 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 let mut minted_batches = VecDeque::new();
1336
1337 let catchup_start = if *resume_upper == [Timestamp::minimum()] {
1340 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 batch_upper
1350 } else {
1351 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 if catchup_start.is_empty() {
1376 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 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 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 return Ok(());
1456 }
1457 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#[derive(Clone, Debug, Serialize, Deserialize)]
1516struct AvroDataFile {
1517 pub data_file: Vec<u8>,
1518 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#[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#[derive(Clone, Debug, Default)]
1595struct BoundedDataFileSet {
1596 pub data_files: Vec<BoundedDataFile>,
1597}
1598
1599fn data_file_location(configured_path: Option<&str>, location: &str) -> String {
1609 if let Some(path) = configured_path {
1616 return path.trim_end_matches('/').to_string();
1617 }
1618
1619 let corrected_location = match location.rsplit_once("/metadata/") {
1623 Some((a, b)) if b.ends_with(".metadata.json") => a,
1624 _ => location,
1625 };
1626 format!("{}/data", corrected_location.trim_end_matches('/'))
1629}
1630
1631fn 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 }
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 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 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 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(¤t_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 let mut stashed_rows: VecDeque<ArcBatch<OrdValBatch<_>>> = VecDeque::new();
1758
1759 let mut in_flight_batches: VecDeque<(
1763 (Antichain<Timestamp>, Antichain<Timestamp>),
1764 Box<dyn IcebergWriter>,
1765 )> = VecDeque::new();
1766
1767 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 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 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 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 !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 let mut staged_messages_since_flush: u64 = 0;
1869
1870 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 if let Some(warner) = pk_warner.as_mut() {
1905 warner.flush();
1906 }
1907 Ok(())
1908 };
1909
1910 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 capset.downgrade(batch_desc.1.clone());
1949 Ok(())
1950 };
1951
1952 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
1982type BatchDescription = (Antichain<Timestamp>, Antichain<Timestamp>);
1984
1985async 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 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 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 break;
2040 };
2041
2042 let Some(rows) = input_batches.front() else {
2043 break;
2045 };
2046
2047 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 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 input_batches.pop_front();
2066 }
2067
2068 }
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 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 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 assert_eq!(
2110 manifest_uri(Some("s3://bucket/tbl/data"), "s3://bucket/tbl"),
2111 "s3://bucket/tbl/data/part-00000.parquet"
2112 );
2113
2114 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 assert_eq!(
2130 manifest_uri(None, "s3://bucket/tbl"),
2131 "s3://bucket/tbl/data/part-00000.parquet"
2132 );
2133
2134 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 assert_eq!(
2146 data_file_location(None, "s3://bucket/tbl/metadata/00001-abc.metadata.json"),
2147 "s3://bucket/tbl/data"
2148 );
2149
2150 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 let result = iceberg_type_overrides(&SqlScalarType::UInt16);
2161 assert_eq!(result.unwrap().0, DataType::Int32);
2162
2163 let result = iceberg_type_overrides(&SqlScalarType::UInt32);
2165 assert_eq!(result.unwrap().0, DataType::Int64);
2166
2167 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 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 let result = iceberg_type_overrides(&SqlScalarType::Interval);
2183 assert_eq!(result.unwrap().0, DataType::LargeUtf8);
2184
2185 let result = iceberg_type_overrides(&SqlScalarType::Uuid);
2187 assert_eq!(result.unwrap().0, DataType::Utf8);
2188
2189 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 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 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 let result = iceberg_type_overrides(&SqlScalarType::Interval);
2229 assert_eq!(result.unwrap().0, DataType::LargeUtf8);
2230
2231 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 assert_eq!(arrow_schema.field(1).data_type(), &DataType::LargeUtf8);
2242
2243 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 #[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 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 #[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 let legacy = with_field_type(&expected, "u", Type::Primitive(PrimitiveType::Fixed(16)));
2311 assert!(is_compatible(&legacy, &expected));
2312
2313 assert!(!is_compatible(&expected, &legacy));
2316
2317 let wrong_type = with_field_type(&expected, "u", Type::Primitive(PrimitiveType::Long));
2319 assert!(!is_compatible(&wrong_type, &expected));
2320
2321 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 #[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 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 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 #[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 assert_eq!(entry_fields[0].name(), "key");
2487 assert_eq!(entry_fields[1].name(), "value");
2488 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 fn frontier(t: Option<u64>) -> Antichain<Timestamp> {
2506 t.map_or_else(Antichain::new, |t| Antichain::from_elem(Timestamp::new(t)))
2507 }
2508
2509 fn span(lower: u64, upper: Option<u64>) -> BatchDescription {
2511 (frontier(Some(lower)), frontier(upper))
2512 }
2513
2514 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 Write(BatchDescription, BatchDescription),
2525 Close(BatchDescription),
2526 }
2527
2528 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 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 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 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 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 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 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
2720fn 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 }
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 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 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 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 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 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}