Repository navigation
fix: match Spark 4.2 date_trunc overflow for SECOND and MILLISECOND #6740
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -41,20 +41,25 @@ pub struct TimestampTruncExpr { | |
| /// `Arc<str>` so it can be cheaply cloned onto Arrow `Timestamp` data types without | ||
| /// reallocating, and parsed once into a `chrono::TimeZone` per batch. | ||
| timezone: Arc<str>, | ||
| /// Whether SECOND and MILLISECOND truncation wraps below the smallest timestamp, as Spark | ||
| /// did before 4.2.0, instead of raising `long overflow` as 4.2.0 does (SPARK-56663). | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Is
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I don't think we can predict the behavior of future Spark versions
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Should future Spark versions inherit the 4.2.0 behavior? Otherwise, several comments might need update again when we support 4.3.0.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes. The serde sets the flag from |
||
| wrap_second_millisecond_overflow: bool, | ||
| } | ||
|
|
||
| impl Hash for TimestampTruncExpr { | ||
| fn hash<H: std::hash::Hasher>(&self, state: &mut H) { | ||
| self.child.hash(state); | ||
| self.format.hash(state); | ||
| self.timezone.hash(state); | ||
| self.wrap_second_millisecond_overflow.hash(state); | ||
| } | ||
| } | ||
| impl PartialEq for TimestampTruncExpr { | ||
| fn eq(&self, other: &Self) -> bool { | ||
| self.child.eq(&other.child) | ||
| && self.format.eq(&other.format) | ||
| && self.timezone.eq(&other.timezone) | ||
| && self.wrap_second_millisecond_overflow == other.wrap_second_millisecond_overflow | ||
| } | ||
| } | ||
|
|
||
|
|
@@ -63,11 +68,13 @@ impl TimestampTruncExpr { | |
| child: Arc<dyn PhysicalExpr>, | ||
| format: Arc<dyn PhysicalExpr>, | ||
| timezone: String, | ||
| wrap_second_millisecond_overflow: bool, | ||
| ) -> Self { | ||
| TimestampTruncExpr { | ||
| child, | ||
| format, | ||
| timezone: Arc::from(timezone), | ||
| wrap_second_millisecond_overflow, | ||
| } | ||
| } | ||
| } | ||
|
|
@@ -100,6 +107,7 @@ impl PhysicalExpr for TimestampTruncExpr { | |
| let format = self.format.evaluate(batch)?; | ||
| let output_type = output_type(×tamp.data_type()); | ||
| let tz = &self.timezone; | ||
| let wrap = self.wrap_second_millisecond_overflow; | ||
| let resolve_tz = |ts: ArrayRef| -> datafusion::common::Result<ArrayRef> { | ||
| // For TimestampNTZ (Timestamp(Microsecond, None)), skip timezone conversion. | ||
| // NTZ values are timezone-independent and truncation should operate directly on the | ||
|
|
@@ -121,15 +129,16 @@ impl PhysicalExpr for TimestampTruncExpr { | |
| }; | ||
| match (timestamp, format) { | ||
| (ColumnarValue::Array(ts), ColumnarValue::Scalar(Utf8(Some(format)))) => { | ||
| let result = timestamp_trunc_dyn(&resolve_tz(ts)?, format)?; | ||
| let result = timestamp_trunc_dyn(&resolve_tz(ts)?, format, wrap)?; | ||
| Ok(ColumnarValue::Array(relabel(result)?)) | ||
| } | ||
| (ColumnarValue::Array(ts), ColumnarValue::Array(formats)) => { | ||
| let result = timestamp_trunc_array_fmt_dyn(&resolve_tz(ts)?, &formats)?; | ||
| Ok(ColumnarValue::Array(relabel(result)?)) | ||
| } | ||
| (ColumnarValue::Scalar(ts_scalar), ColumnarValue::Scalar(Utf8(Some(format)))) => { | ||
| let result = timestamp_trunc_dyn(&resolve_tz(ts_scalar.to_array()?)?, format)?; | ||
| let result = | ||
| timestamp_trunc_dyn(&resolve_tz(ts_scalar.to_array()?)?, format, wrap)?; | ||
| let scalar = ScalarValue::try_from_array(&relabel(result)?, 0)?; | ||
| Ok(ColumnarValue::Scalar(scalar)) | ||
| } | ||
|
|
@@ -153,6 +162,7 @@ impl PhysicalExpr for TimestampTruncExpr { | |
| Arc::clone(&children[0]), | ||
| Arc::clone(&self.format), | ||
| self.timezone.to_string(), | ||
| self.wrap_second_millisecond_overflow, | ||
| ))) | ||
| } | ||
| } | ||
|
|
@@ -195,6 +205,7 @@ mod tests { | |
| Arc::new(Column::new("ts", 0)), | ||
| Arc::new(Literal::new(Utf8(Some("HOUR".to_string())))), | ||
| "Asia/Kolkata".to_string(), | ||
| false, | ||
| ); | ||
| let declared = expr.data_type(&schema).unwrap(); | ||
| let ColumnarValue::Array(result) = expr.evaluate(&batch).unwrap() else { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This is to track the Spark version Comet is running against?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Yes
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
If so, why it's not named as
spark_420_plus?There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think it is better to name the behavior rather than a Spark version. Some companies maintain forks of Spark and backport fixes from newer OSS Spark versions. These companies can leverate the shim layer to say which of their versions have specific behavior. I have also seen in the past (although rarely) that there are functional differences between different minor releases.