Skip to content

Commit efa47cf

Browse files
fix(delta): strip system columns and normalize timestamps for writes
Exclude _timestamp/_updating_meta from Delta catalog and Parquet schema, cast user columns to Delta-compatible types, and support multiple timestamp units in watermark max extraction. Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 76d0204 commit efa47cf

4 files changed

Lines changed: 230 additions & 24 deletions

File tree

src/streaming_runtime_operators/src/factory/connector/delta.rs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,9 @@ use crate::factory::connector::sink_props_codec::{
2424
};
2525
use crate::factory::global::Registry;
2626
use crate::factory::operator_constructor::OperatorConstructor;
27-
use crate::operators::sink::delta::{DeltaFormat, DeltaSinkOperator};
27+
use crate::operators::sink::delta::{
28+
DeltaFormat, DeltaSinkOperator, strip_streaming_system_columns_arc,
29+
};
2830
use crate::operators::sink::filesystem::compression_from_str;
2931
use crate::sql::common::FsSchema;
3032
use crate::sql::common::constants::connection_format_value;
@@ -47,7 +49,7 @@ impl OperatorConstructor for DeltaSinkDispatcher {
4749
.map(|fs| FsSchema::try_from(fs.clone()))
4850
.transpose()
4951
.map_err(|e| anyhow::anyhow!("invalid fs_schema for delta sink: {e}"))?
50-
.map(|fs| fs.schema);
52+
.and_then(|fs| strip_streaming_system_columns_arc(fs.schema));
5153

5254
let format = props
5355
.get(opt::FORMAT)

src/streaming_runtime_operators/src/operators/sink/delta/delta_commit.rs

Lines changed: 161 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,10 @@ use std::collections::HashMap;
1414
use std::sync::Arc;
1515
use std::time::{SystemTime, UNIX_EPOCH};
1616

17+
use arrow::compute::cast;
1718
use arrow_array::RecordBatch;
1819
use arrow_ipc::writer::StreamWriter;
19-
use arrow_schema::Schema as ArrowSchema;
20+
use arrow_schema::{DataType, Field, FieldRef, Schema as ArrowSchema, TimeUnit};
2021
use deltalake::errors::DeltaTableError;
2122
use deltalake::kernel::engine::arrow_conversion::TryIntoKernel as _;
2223
use deltalake::kernel::schema::cast::normalize_for_delta;
@@ -27,8 +28,42 @@ use deltalake::{DeltaTable, open_table_with_storage_options};
2728
use tracing::{info, instrument};
2829
use url::Url;
2930

31+
use crate::sql::common::{TIMESTAMP_FIELD, UPDATING_META_FIELD};
32+
3033
use super::DeltaSinkError;
3134

35+
/// Streaming-internal columns that must not be persisted to external sinks.
36+
pub fn is_streaming_system_column(name: &str) -> bool {
37+
name == TIMESTAMP_FIELD || name == UPDATING_META_FIELD
38+
}
39+
40+
/// Remove `_timestamp` / `_updating_meta` from a schema (e.g. connector `fs_schema` may inject them).
41+
pub fn strip_streaming_system_columns(schema: &ArrowSchema) -> ArrowSchema {
42+
let fields: Vec<FieldRef> = schema
43+
.fields()
44+
.iter()
45+
.filter(|f| !is_streaming_system_column(f.name()))
46+
.cloned()
47+
.collect();
48+
ArrowSchema::new(fields)
49+
}
50+
51+
pub fn strip_streaming_system_columns_arc(schema: Arc<ArrowSchema>) -> Option<Arc<ArrowSchema>> {
52+
let stripped = strip_streaming_system_columns(schema.as_ref());
53+
if stripped.fields().is_empty() {
54+
return None;
55+
}
56+
let had_system = schema
57+
.fields()
58+
.iter()
59+
.any(|f| is_streaming_system_column(f.name()));
60+
if had_system {
61+
Some(Arc::new(stripped))
62+
} else {
63+
Some(schema)
64+
}
65+
}
66+
3267
pub struct UncommittedDataFile {
3368
pub path: String,
3469
pub size_bytes: u64,
@@ -42,6 +77,8 @@ pub struct DeltaTableCommitter {
4277
table: Option<DeltaTable>,
4378
/// Precomputed from catalog `fs_schema` at startup, or lazily from the first batch.
4479
delta_columns: Option<Vec<StructField>>,
80+
/// Arrow 55 schema for Parquet writes (timestamps normalized to microsecond).
81+
write_schema: Option<Arc<ArrowSchema>>,
4582
}
4683

4784
impl DeltaTableCommitter {
@@ -50,7 +87,11 @@ impl DeltaTableCommitter {
5087
storage_options: HashMap<String, String>,
5188
catalog_schema: Option<Arc<ArrowSchema>>,
5289
) -> Result<Self, DeltaSinkError> {
53-
let delta_columns = catalog_schema
90+
let user_schema = catalog_schema.and_then(strip_streaming_system_columns_arc);
91+
let write_schema = user_schema
92+
.as_ref()
93+
.map(|s| Arc::new(normalize_arrow_schema_for_delta(s)));
94+
let delta_columns = user_schema
5495
.as_deref()
5596
.map(arrow_schema_to_delta_columns)
5697
.transpose()?;
@@ -61,13 +102,30 @@ impl DeltaTableCommitter {
61102
uncommitted: Vec::new(),
62103
table: None,
63104
delta_columns,
105+
write_schema,
64106
})
65107
}
66108

67-
/// Fallback when catalog schema is absent: derive columns from the first flushed batch.
109+
pub fn write_schema(&self) -> Option<Arc<ArrowSchema>> {
110+
self.write_schema.clone()
111+
}
112+
113+
/// Fallback when catalog schema is absent: derive user columns from the first flushed batch.
68114
pub fn update_schema(&mut self, schema: Arc<ArrowSchema>) -> Result<(), DeltaSinkError> {
115+
let user_schema = strip_streaming_system_columns_arc(schema).ok_or_else(|| {
116+
DeltaSinkError::CommitterFailed(
117+
"cannot derive delta table schema: no user columns after removing streaming \
118+
system columns (_timestamp, _updating_meta)"
119+
.into(),
120+
)
121+
})?;
69122
if self.delta_columns.is_none() {
70-
self.delta_columns = Some(arrow_schema_to_delta_columns(&schema)?);
123+
self.delta_columns = Some(arrow_schema_to_delta_columns(user_schema.as_ref())?);
124+
}
125+
if self.write_schema.is_none() {
126+
self.write_schema = Some(Arc::new(normalize_arrow_schema_for_delta(
127+
user_schema.as_ref(),
128+
)));
71129
}
72130
Ok(())
73131
}
@@ -231,6 +289,105 @@ impl DeltaTableCommitter {
231289
}
232290
}
233291

292+
/// Normalize Arrow 55 schema for Delta Parquet writes (align with deltalake `normalize_for_delta`).
293+
pub fn normalize_arrow_schema_for_delta(schema: &ArrowSchema) -> ArrowSchema {
294+
let fields: Vec<FieldRef> = schema
295+
.fields()
296+
.iter()
297+
.map(|f| Arc::new(normalize_field_for_delta(f.as_ref())))
298+
.collect();
299+
ArrowSchema::new(fields)
300+
}
301+
302+
fn normalize_field_for_delta(field: &Field) -> Field {
303+
let data_type = normalize_datatype_for_delta(field.data_type());
304+
if data_type == *field.data_type() {
305+
field.clone()
306+
} else {
307+
field.clone().with_data_type(data_type)
308+
}
309+
}
310+
311+
fn normalize_datatype_for_delta(dt: &DataType) -> DataType {
312+
match dt {
313+
DataType::Date64 => DataType::Date32,
314+
DataType::Timestamp(TimeUnit::Second, tz)
315+
| DataType::Timestamp(TimeUnit::Millisecond, tz)
316+
| DataType::Timestamp(TimeUnit::Nanosecond, tz) => {
317+
DataType::Timestamp(TimeUnit::Microsecond, tz.clone())
318+
}
319+
DataType::Struct(fields) => {
320+
let normalized: Vec<FieldRef> = fields
321+
.iter()
322+
.map(|f| Arc::new(normalize_field_for_delta(f.as_ref())))
323+
.collect();
324+
DataType::Struct(normalized.into())
325+
}
326+
DataType::List(inner) => {
327+
DataType::List(Arc::new(normalize_field_for_delta(inner.as_ref())))
328+
}
329+
DataType::LargeList(inner) => {
330+
DataType::LargeList(Arc::new(normalize_field_for_delta(inner.as_ref())))
331+
}
332+
DataType::FixedSizeList(inner, len) => {
333+
DataType::FixedSizeList(Arc::new(normalize_field_for_delta(inner.as_ref())), *len)
334+
}
335+
DataType::Map(entries, sorted) => DataType::Map(
336+
Arc::new(normalize_field_for_delta(entries.as_ref())),
337+
*sorted,
338+
),
339+
_ => dt.clone(),
340+
}
341+
}
342+
343+
/// Cast record batches so on-disk Parquet matches the Delta table schema.
344+
pub fn cast_batches_for_delta_write(
345+
batches: &[RecordBatch],
346+
target_schema: &ArrowSchema,
347+
) -> Result<Vec<RecordBatch>, DeltaSinkError> {
348+
let target = Arc::new(target_schema.clone());
349+
batches
350+
.iter()
351+
.map(|batch| cast_batch_for_delta_write(batch, &target))
352+
.collect()
353+
}
354+
355+
fn cast_batch_for_delta_write(
356+
batch: &RecordBatch,
357+
target_schema: &Arc<ArrowSchema>,
358+
) -> Result<RecordBatch, DeltaSinkError> {
359+
if batch.schema().as_ref() == target_schema.as_ref() {
360+
return Ok(batch.clone());
361+
}
362+
363+
let mut columns = Vec::with_capacity(target_schema.fields().len());
364+
for field in target_schema.fields() {
365+
let col = batch.column_by_name(field.name()).ok_or_else(|| {
366+
DeltaSinkError::CommitterFailed(format!(
367+
"batch missing column '{}' required by delta write schema",
368+
field.name()
369+
))
370+
})?;
371+
let casted = if col.data_type() == field.data_type() {
372+
col.clone()
373+
} else {
374+
cast(col, field.data_type()).map_err(|e| {
375+
DeltaSinkError::CommitterFailed(format!(
376+
"failed to cast column '{}' from {:?} to {:?}: {e}",
377+
field.name(),
378+
col.data_type(),
379+
field.data_type()
380+
))
381+
})?
382+
};
383+
columns.push(casted);
384+
}
385+
386+
RecordBatch::try_new(target_schema.clone(), columns).map_err(|e| {
387+
DeltaSinkError::CommitterFailed(format!("failed to build delta write batch: {e}"))
388+
})
389+
}
390+
234391
/// Bridge arrow 55 (runtime) schema to deltalake kernel schema via IPC.
235392
fn arrow_schema_to_delta_columns(schema: &ArrowSchema) -> Result<Vec<StructField>, DeltaSinkError> {
236393
let empty = RecordBatch::new_empty(Arc::new(schema.clone()));

src/streaming_runtime_operators/src/operators/sink/delta/mod.rs

Lines changed: 34 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,8 @@
1212

1313
mod delta_commit;
1414

15+
pub use delta_commit::strip_streaming_system_columns_arc;
16+
1517
use std::collections::HashMap;
1618
use std::path::PathBuf;
1719
use std::sync::Arc;
@@ -21,7 +23,8 @@ use arrow_schema::Schema as ArrowSchema;
2123
use async_trait::async_trait;
2224
use bytes::Bytes;
2325
use delta_commit::{
24-
DeltaTableCommitter, UncommittedDataFile, build_delta_storage_options, resolve_delta_table_uri,
26+
DeltaTableCommitter, UncommittedDataFile, build_delta_storage_options,
27+
cast_batches_for_delta_write, resolve_delta_table_uri,
2528
};
2629
use object_store::aws::AmazonS3Builder;
2730
use object_store::path::Path as ObjectStorePath;
@@ -94,6 +97,8 @@ pub struct DeltaSinkOperator {
9497
s3_bucket: Option<String>,
9598
sink_path: String,
9699
catalog_schema: Option<Arc<ArrowSchema>>,
100+
/// Normalized schema for Parquet files (microsecond timestamps, etc.).
101+
parquet_write_schema: Option<Arc<ArrowSchema>>,
97102
}
98103

99104
impl DeltaSinkOperator {
@@ -171,7 +176,8 @@ impl DeltaSinkOperator {
171176
storage_options,
172177
s3_bucket,
173178
sink_path: path,
174-
catalog_schema,
179+
catalog_schema: catalog_schema.and_then(strip_streaming_system_columns_arc),
180+
parquet_write_schema: None,
175181
})
176182
}
177183

@@ -186,23 +192,35 @@ impl DeltaSinkOperator {
186192
return Ok(());
187193
}
188194

189-
let fallback_schema = self
190-
.catalog_schema
191-
.is_none()
192-
.then(|| self.pending[0].schema());
195+
let fallback_schema = self.catalog_schema.is_none().then(|| {
196+
strip_streaming_system_columns_arc(self.pending[0].schema())
197+
.expect("batch has no user columns after removing streaming system columns")
198+
});
193199
let batches = std::mem::take(&mut self.pending);
194200
self.pending_bytes = 0;
195201

196202
let record_count: u64 = batches.iter().map(|b| b.num_rows() as u64).sum();
197203
let format = self.format;
198204
let compression = self.parquet_compression;
205+
let parquet_write_schema = self.parquet_write_schema.clone();
199206

200-
let bytes = tokio::task::spawn_blocking(move || match format {
201-
DeltaFormat::Csv => FormatEncoder::encode_csv(&batches),
202-
DeltaFormat::Parquet => FormatEncoder::encode_parquet(&batches, compression),
203-
DeltaFormat::JsonL => FormatEncoder::encode_jsonl(&batches),
204-
DeltaFormat::Avro => FormatEncoder::encode_avro(&batches),
205-
DeltaFormat::Orc => FormatEncoder::encode_orc(&batches),
207+
let bytes = tokio::task::spawn_blocking(move || {
208+
let batches = if format == DeltaFormat::Parquet {
209+
if let Some(ref schema) = parquet_write_schema {
210+
cast_batches_for_delta_write(&batches, schema)?
211+
} else {
212+
batches
213+
}
214+
} else {
215+
batches
216+
};
217+
match format {
218+
DeltaFormat::Csv => FormatEncoder::encode_csv(&batches),
219+
DeltaFormat::Parquet => FormatEncoder::encode_parquet(&batches, compression),
220+
DeltaFormat::JsonL => FormatEncoder::encode_jsonl(&batches),
221+
DeltaFormat::Avro => FormatEncoder::encode_avro(&batches),
222+
DeltaFormat::Orc => FormatEncoder::encode_orc(&batches),
223+
}
206224
})
207225
.await
208226
.map_err(|e| DeltaSinkError::SerializationPanic(e.to_string()))?
@@ -302,11 +320,13 @@ impl Operator for DeltaSinkOperator {
302320
self.table_uri = Some(table_uri.clone());
303321

304322
if self.format == DeltaFormat::Parquet {
305-
self.committer = Some(DeltaTableCommitter::try_new(
323+
let committer = DeltaTableCommitter::try_new(
306324
table_uri,
307325
self.storage_options.clone(),
308326
self.catalog_schema.clone(),
309-
)?);
327+
)?;
328+
self.parquet_write_schema = committer.write_schema();
329+
self.committer = Some(committer);
310330
} else {
311331
warn!(
312332
format = ?self.format,

src/streaming_runtime_operators/src/operators/watermark/watermark_generator.rs

Lines changed: 31 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,12 @@
1313
use anyhow::{Result, anyhow};
1414
use arrow::compute::kernels::aggregate;
1515
use arrow_array::cast::AsArray;
16-
use arrow_array::types::TimestampNanosecondType;
16+
use arrow_array::types::{
17+
TimestampMicrosecondType, TimestampMillisecondType, TimestampNanosecondType,
18+
TimestampSecondType,
19+
};
1720
use arrow_array::{RecordBatch, TimestampNanosecondArray};
21+
use arrow_schema::{DataType, TimeUnit};
1822
use bincode::{Decode, Encode};
1923
use datafusion::physical_expr::PhysicalExpr;
2024
use datafusion_proto::physical_plan::DefaultPhysicalExtensionCodec;
@@ -78,9 +82,32 @@ impl WatermarkGeneratorOperator {
7882

7983
fn extract_max_timestamp(&self, batch: &RecordBatch) -> Option<SystemTime> {
8084
let ts_column = batch.column(self.timestamp_index);
81-
let arr = ts_column.as_primitive::<TimestampNanosecondType>();
82-
let max_ts = aggregate::max(arr)?;
83-
Some(from_nanos(max_ts as u128))
85+
match ts_column.data_type() {
86+
DataType::Timestamp(TimeUnit::Nanosecond, _) => {
87+
let arr = ts_column.as_primitive::<TimestampNanosecondType>();
88+
aggregate::max(arr).map(|v| from_nanos(v as u128))
89+
}
90+
DataType::Timestamp(TimeUnit::Microsecond, _) => {
91+
let arr = ts_column.as_primitive::<TimestampMicrosecondType>();
92+
aggregate::max(arr).map(|v| from_nanos((v as u128) * 1_000))
93+
}
94+
DataType::Timestamp(TimeUnit::Millisecond, _) => {
95+
let arr = ts_column.as_primitive::<TimestampMillisecondType>();
96+
aggregate::max(arr).map(|v| from_nanos((v as u128) * 1_000_000))
97+
}
98+
DataType::Timestamp(TimeUnit::Second, _) => {
99+
let arr = ts_column.as_primitive::<TimestampSecondType>();
100+
aggregate::max(arr).map(|v| from_nanos((v as u128) * 1_000_000_000))
101+
}
102+
other => {
103+
debug!(
104+
?other,
105+
index = self.timestamp_index,
106+
"skip max timestamp: column is not a timestamp array"
107+
);
108+
None
109+
}
110+
}
84111
}
85112

86113
fn evaluate_watermark(&self, batch: &RecordBatch) -> Result<SystemTime> {

0 commit comments

Comments
 (0)