From 92ef77fd67f7c957c58c9763a28ae705e010bd23 Mon Sep 17 00:00:00 2001 From: Ashutosh Date: Tue, 11 Aug 2026 14:49:17 +0530 Subject: [PATCH 1/2] fix: iceberg dummy partition for every data --- .../sinks/iceberg_sink/src/router/mod.rs | 202 +++++++++--------- 1 file changed, 104 insertions(+), 98 deletions(-) diff --git a/core/connectors/sinks/iceberg_sink/src/router/mod.rs b/core/connectors/sinks/iceberg_sink/src/router/mod.rs index 3f43b94013..9ccd8e66ac 100644 --- a/core/connectors/sinks/iceberg_sink/src/router/mod.rs +++ b/core/connectors/sinks/iceberg_sink/src/router/mod.rs @@ -21,14 +21,14 @@ use arrow_json::ReaderBuilder; use async_trait::async_trait; use iceberg::TableIdent; use iceberg::arrow::schema_to_arrow_schema; -use iceberg::spec::{ - Literal, PartitionKey, PartitionSpec, PrimitiveLiteral, PrimitiveType, Struct, StructType, -}; +use iceberg::arrow::RecordBatchPartitionSplitter; use iceberg::table::Table; use iceberg::transaction::{ApplyTransactionAction, Transaction}; use iceberg::writer::base_writer::data_file_writer::DataFileWriterBuilder; use iceberg::writer::file_writer::ParquetWriterBuilder; use iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder; +use iceberg::writer::partitioning::fanout_writer::FanoutWriter; +use iceberg::writer::partitioning::PartitioningWriter; use iceberg::writer::{IcebergWriter, IcebergWriterBuilder}; use iceberg::{ Catalog, @@ -67,49 +67,6 @@ async fn table_exists(route_field_val: &str, catalog: &dyn Catalog) -> Option Result { - match pt { - PrimitiveType::Boolean => Ok(PrimitiveLiteral::Boolean(false)), - PrimitiveType::Int => Ok(PrimitiveLiteral::Int(0)), - PrimitiveType::Long => Ok(PrimitiveLiteral::Long(0)), - PrimitiveType::Decimal { .. } => Ok(PrimitiveLiteral::Int128(0)), - PrimitiveType::Date => Ok(PrimitiveLiteral::Int(0)), // e.g. days since epoch - PrimitiveType::Time => Ok(PrimitiveLiteral::Long(0)), // microseconds since midnight - PrimitiveType::Timestamp => Ok(PrimitiveLiteral::Long(0)), // microseconds since epoch - PrimitiveType::Timestamptz => Ok(PrimitiveLiteral::Long(0)), - PrimitiveType::TimestampNs => Ok(PrimitiveLiteral::Long(0)), - PrimitiveType::TimestamptzNs => Ok(PrimitiveLiteral::Long(0)), - PrimitiveType::String => Ok(PrimitiveLiteral::String(String::new())), - PrimitiveType::Uuid => Ok(PrimitiveLiteral::Binary(vec![0; 16])), - PrimitiveType::Fixed(len) => Ok(PrimitiveLiteral::Binary(vec![0; *len as usize])), - PrimitiveType::Binary => Ok(PrimitiveLiteral::Binary(Vec::new())), - _ => { - error!("Partition type not supported"); - Err(Error::InvalidConfig) - } - } -} - -fn get_partition_type_value(default_partition_type: &StructType) -> Result, Error> { - let mut fields: Vec> = Vec::new(); - - if default_partition_type.fields().is_empty() { - return Ok(None); - }; - - for field in default_partition_type.fields() { - let field_type = field.field_type.as_primitive_type().ok_or_else(|| { - error!("The partition type of the configured iceberg table is not a primitive type"); - Error::InvalidConfig - })?; - - let value = Some(Literal::Primitive(primitive_type_to_literal(field_type)?)); - - fields.push(value); - } - Ok(Some(Struct::from_iter(fields))) -} - async fn write_data( messages: &[Payload], table: &Table, @@ -145,27 +102,6 @@ async fn write_data( let data_file_writer_builder = DataFileWriterBuilder::new(rolling_file_writer_builder); - let partition_spec = PartitionSpec::builder(table.current_schema_ref()); - - let partition_type = get_partition_type_value(table.metadata().default_partition_type())?; - - let mut writer = data_file_writer_builder - .build(match partition_type { - None => None, - Some(p_type) => Some(PartitionKey::new( - partition_spec - .build() - .map_err(|err| Error::InitError(err.to_string()))?, - table.current_schema_ref(), - p_type, - )), - }) - .await - .map_err(|err| { - error!("Error while constructing data file writer: {}", err); - Error::InitError(err.to_string()) - })?; - let msgs: Vec<&simd_json::OwnedValue> = messages .iter() .filter_map(|payload| match payload { @@ -201,42 +137,112 @@ async fn write_data( Error::InitError(err.to_string()) })?; - let write_result: Result<(), Error> = async { - for batch in reader { - let batch_data = batch.map_err(|err| { - let chain = format_error_chain(&err); - error!("Error while getting record batch: {}", chain); - Error::InvalidRecordValue(chain) - })?; - writer.write(batch_data).await.map_err(|err| { - let chain = format_error_chain(&err); - error!("Error while writing record batch: {}", chain); - Error::WriteFailure(chain) + let partition_spec = table.metadata().default_partition_spec(); + + let data_files = if partition_spec.is_unpartitioned() { + let mut writer = data_file_writer_builder + .build(None) + .await + .map_err(|err| { + error!("Error while constructing data file writer: {}", err); + Error::InitError(err.to_string()) })?; + + let write_result: Result<(), Error> = async { + for batch in reader { + let batch_data = batch.map_err(|err| { + let chain = format_error_chain(&err); + error!("Error while getting record batch: {}", chain); + Error::InvalidRecordValue(chain) + })?; + writer.write(batch_data).await.map_err(|err| { + let chain = format_error_chain(&err); + error!("Error while writing record batch: {}", chain); + Error::WriteFailure(chain) + })?; + } + Ok(()) } - Ok(()) - } - .await; + .await; - if let Err(e) = &write_result { - error!( - "Batch loop failed ({}), closing writer to release resources", - e - ); - if let Err(close_err) = writer.close().await { - error!("Failed to close writer after batch error: {}", close_err); + if let Err(e) = &write_result { + error!( + "Batch loop failed ({}), closing writer to release resources", + e + ); + if let Err(close_err) = writer.close().await { + error!("Failed to close writer after batch error: {}", close_err); + } + return Err(write_result.unwrap_err()); } - return Err(write_result.unwrap_err()); - } - let data_files = writer.close().await.map_err(|err| { - let chain = format_error_chain(&err); - error!( - "Error while writing data records to Parquet file: {}", - chain - ); - Error::WriteFailure(chain) - })?; + writer.close().await.map_err(|err| { + let chain = format_error_chain(&err); + error!( + "Error while writing data records to Parquet file: {}", + chain + ); + Error::WriteFailure(chain) + })? + } else { + let splitter = RecordBatchPartitionSplitter::try_new_with_computed_values( + table.metadata().current_schema().clone(), + partition_spec.clone(), + ) + .map_err(|err| { + error!("Failed to create partition splitter: {}", err); + Error::InitError(err.to_string()) + })?; + + let mut fanout_writer = FanoutWriter::new(data_file_writer_builder); + + let write_result: Result<(), Error> = async { + for batch in reader { + let batch_data = batch.map_err(|err| { + let chain = format_error_chain(&err); + error!("Error while getting record batch: {}", chain); + Error::InvalidRecordValue(chain) + })?; + let partitioned_batches = splitter.split(&batch_data).map_err(|err| { + let chain = format_error_chain(&err); + error!("Error while splitting batch by partition: {}", chain); + Error::InvalidRecordValue(chain) + })?; + for (partition_key, partition_batch) in partitioned_batches { + fanout_writer + .write(partition_key, partition_batch) + .await + .map_err(|err| { + let chain = format_error_chain(&err); + error!("Error while writing record batch: {}", chain); + Error::WriteFailure(chain) + })?; + } + } + Ok(()) + } + .await; + + if let Err(e) = &write_result { + error!( + "Batch loop failed ({}), closing writer to release resources", + e + ); + if let Err(close_err) = fanout_writer.close().await { + error!("Failed to close writer after batch error: {}", close_err); + } + return Err(write_result.unwrap_err()); + } + + fanout_writer.close().await.map_err(|err| { + let chain = format_error_chain(&err); + error!( + "Error while writing data records to Parquet file: {}", + chain + ); + Error::WriteFailure(chain) + })? + }; let table_commit = Transaction::new(table); From dc10410c5a240d282120be52e7e8a0624150c548 Mon Sep 17 00:00:00 2001 From: Ashutosh Date: Tue, 11 Aug 2026 15:01:18 +0530 Subject: [PATCH 2/2] chore: cargo fmt --- .../sinks/iceberg_sink/src/router/mod.rs | 15 ++++++--------- 1 file changed, 6 insertions(+), 9 deletions(-) diff --git a/core/connectors/sinks/iceberg_sink/src/router/mod.rs b/core/connectors/sinks/iceberg_sink/src/router/mod.rs index 9ccd8e66ac..0ec10a92b2 100644 --- a/core/connectors/sinks/iceberg_sink/src/router/mod.rs +++ b/core/connectors/sinks/iceberg_sink/src/router/mod.rs @@ -20,15 +20,15 @@ use crate::slice_user_table; use arrow_json::ReaderBuilder; use async_trait::async_trait; use iceberg::TableIdent; -use iceberg::arrow::schema_to_arrow_schema; use iceberg::arrow::RecordBatchPartitionSplitter; +use iceberg::arrow::schema_to_arrow_schema; use iceberg::table::Table; use iceberg::transaction::{ApplyTransactionAction, Transaction}; use iceberg::writer::base_writer::data_file_writer::DataFileWriterBuilder; use iceberg::writer::file_writer::ParquetWriterBuilder; use iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder; -use iceberg::writer::partitioning::fanout_writer::FanoutWriter; use iceberg::writer::partitioning::PartitioningWriter; +use iceberg::writer::partitioning::fanout_writer::FanoutWriter; use iceberg::writer::{IcebergWriter, IcebergWriterBuilder}; use iceberg::{ Catalog, @@ -140,13 +140,10 @@ async fn write_data( let partition_spec = table.metadata().default_partition_spec(); let data_files = if partition_spec.is_unpartitioned() { - let mut writer = data_file_writer_builder - .build(None) - .await - .map_err(|err| { - error!("Error while constructing data file writer: {}", err); - Error::InitError(err.to_string()) - })?; + let mut writer = data_file_writer_builder.build(None).await.map_err(|err| { + error!("Error while constructing data file writer: {}", err); + Error::InitError(err.to_string()) + })?; let write_result: Result<(), Error> = async { for batch in reader {