Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 66 additions & 1 deletion crates/paimon/src/variant.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2369,7 +2369,7 @@ fn variant_shredding_row_type(
fields.push(DataField::new(
2,
VARIANT_TYPED_VALUE_FIELD_NAME.to_string(),
data_type.clone(),
shredding_typed_value_type(data_type)?,
));
}
other => return invalid_variant_shredding_schema(format!("{other:?}")),
Expand All @@ -2378,6 +2378,31 @@ fn variant_shredding_row_type(
Ok(RowType::new(fields))
}

/// The physical `typed_value` type for a shredded Variant scalar.
///
/// The Variant binary stores timestamps as microseconds (encoding type codes 12
/// and 13 are `TIMESTAMP(MICROS)` / `TIMESTAMP_NTZ(MICROS)`), and the shredding
/// spec requires the shredded `typed_value` to use microseconds too. A
/// configured shredding schema may nonetheless declare another precision (e.g.
/// `TIMESTAMP(3)`); if that flowed through verbatim, the writer would place the
/// microsecond value into a millisecond/second/nanosecond column and annotate it
/// as such, so any precision-aware reader (Spark, DuckDB, Arrow, Java
/// `Timestamp.fromMicros`) would read it off by a factor of 1000. Pin the
/// shredded timestamp precision to microseconds; other scalars pass through.
fn shredding_typed_value_type(data_type: &DataType) -> Result<DataType> {
const MICROS_PRECISION: u32 = 6;
Ok(match data_type {
DataType::Timestamp(_) => DataType::Timestamp(TimestampType::with_nullable(
data_type.is_nullable(),
MICROS_PRECISION,
)?),
DataType::LocalZonedTimestamp(_) => DataType::LocalZonedTimestamp(
LocalZonedTimestampType::with_nullable(data_type.is_nullable(), MICROS_PRECISION)?,
),
other => other.clone(),
})
}

fn variant_binary_type(nullable: bool) -> Result<DataType> {
Ok(DataType::VarBinary(VarBinaryType::try_new(
nullable,
Expand Down Expand Up @@ -3276,6 +3301,46 @@ mod tests {
assert!(format!("{err:?}").contains("Invalid variant shredding schema"));
}

#[test]
fn shredded_timestamp_typed_value_is_pinned_to_micros() {
// A configured shredding schema may declare TIMESTAMP(3), but the
// Variant binary stores timestamps as microseconds, so the shredded
// typed_value must be microseconds (precision 6). Otherwise the micros
// value is written into a millisecond-annotated column and read back
// x1000 wrong by any precision-aware engine.
let cases = [
DataType::Timestamp(TimestampType::new(3).unwrap()),
DataType::LocalZonedTimestamp(LocalZonedTimestampType::new(0).unwrap()),
DataType::Timestamp(TimestampType::new(9).unwrap()),
];
for configured in cases {
let DataType::Row(physical) = variant_shredding_type(&configured).unwrap() else {
panic!("expected row shredding type");
};
let typed_value = physical.fields()[2].data_type();
let precision = match typed_value {
DataType::Timestamp(t) => t.precision(),
DataType::LocalZonedTimestamp(t) => t.precision(),
other => panic!("expected timestamp typed_value, got {other:?}"),
};
assert_eq!(
precision, 6,
"shredded timestamp typed_value must be micros, got precision {precision}"
);
}

// Non-timestamp scalars are left untouched.
let DataType::Row(physical) =
variant_shredding_type(&DataType::BigInt(BigIntType::new())).unwrap()
else {
panic!("expected row shredding type");
};
assert!(matches!(
physical.fields()[2].data_type(),
DataType::BigInt(_)
));
}

fn timestamp_variant(value: i64, ntz: bool) -> GenericVariant {
let mut builder = VariantBuilder::new();
if ntz {
Expand Down
Loading