Skip to content

Commit d01ec12

Browse files
committed
DataFrameValueWriter serialized a Spark DateType via java.sql.Date.getTime(),
which is midnight in the JVM default timezone. On a non UTC JVM this shifted the stored value by the timezone offset and could land on a different calendar day. Convert through the calendar date at UTC so the result is timezone independent. Applies to the sql-30, sql-35, and sql-40 modules. Resolves opensearch-project#797 Signed-off-by: Sotaro Hikita <bering1814@gmail.com>
1 parent 408dca7 commit d01ec12

7 files changed

Lines changed: 57 additions & 3 deletions

File tree

CHANGELOG.md

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,9 @@ Inspired from [Keep a Changelog](https://keepachangelog.com/en/1.0.0/)
55
### Added
66
- Add system property `opensearch.hadoop.version.check.skip` to bypass multiple JAR version detection ([#753](https://github.com/opensearch-project/opensearch-hadoop/pull/753))
77

8+
### Fixed
9+
- Write Spark `DateType` values as UTC start of day so the stored date does not shift by the JVM timezone offset on non UTC JVMs ([#797](https://github.com/opensearch-project/opensearch-hadoop/issues/797))
10+
811
### Dependencies
912
- Bumps `commons-logging:commons-logging` from 1.3.5 to 1.3.6
1013
- Bumps `com.fasterxml.jackson.core:jackson-databind` from 2.21.1 to 2.21.3

spark/sql-30/src/main/scala/org/opensearch/spark/sql/DataFrameValueWriter.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -269,7 +269,7 @@ class DataFrameValueWriter(writeUnknownTypes: Boolean = false) extends Filtering
269269
case DoubleType => generator.writeNumber(value.asInstanceOf[Double])
270270
case FloatType => generator.writeNumber(value.asInstanceOf[Float])
271271
case TimestampType => generator.writeNumber(value.asInstanceOf[Timestamp].getTime())
272-
case DateType => generator.writeNumber(value.asInstanceOf[Date].getTime())
272+
case DateType => generator.writeNumber(value.asInstanceOf[Date].toLocalDate.atStartOfDay(java.time.ZoneOffset.UTC).toInstant.toEpochMilli)
273273
case StringType => generator.writeString(value.toString)
274274
case _ => {
275275
val className = schema.getClass().getName()

spark/sql-30/src/test/scala/org/opensearch/spark/sql/DataFrameValueWriterTest.scala

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,10 +30,13 @@
3030
package org.opensearch.spark.sql
3131

3232
import java.io.ByteArrayOutputStream
33+
import java.sql.Date
34+
import java.util.TimeZone
3335

3436
import org.apache.spark.sql.Row
3537
import org.apache.spark.sql.catalyst.ScalaReflection
3638
import org.apache.spark.sql.types.ArrayType
39+
import org.apache.spark.sql.types.DateType
3740
import org.apache.spark.sql.types.IntegerType
3841
import org.apache.spark.sql.types.MapType
3942
import org.apache.spark.sql.types.StringType
@@ -165,4 +168,18 @@ class DataFrameValueWriterTest {
165168
}
166169
}
167170

171+
@Test
172+
def testWriteDateUsesUtcStartOfDay(): Unit = {
173+
val defaultTimeZone = TimeZone.getDefault
174+
try {
175+
TimeZone.setDefault(TimeZone.getTimeZone("Asia/Tokyo"))
176+
val schema = StructType(Seq(StructField("d", DateType)))
177+
val row = Row(Date.valueOf("2023-07-22"))
178+
val serialized = serialize(row, schema)
179+
assertTrue(serialized.contains(""""d":1689984000000"""))
180+
} finally {
181+
TimeZone.setDefault(defaultTimeZone)
182+
}
183+
}
184+
168185
}

spark/sql-35/src/main/scala/org/opensearch/spark/sql/DataFrameValueWriter.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -269,7 +269,7 @@ class DataFrameValueWriter(writeUnknownTypes: Boolean = false) extends Filtering
269269
case DoubleType => generator.writeNumber(value.asInstanceOf[Double])
270270
case FloatType => generator.writeNumber(value.asInstanceOf[Float])
271271
case TimestampType => generator.writeNumber(value.asInstanceOf[Timestamp].getTime())
272-
case DateType => generator.writeNumber(value.asInstanceOf[Date].getTime())
272+
case DateType => generator.writeNumber(value.asInstanceOf[Date].toLocalDate.atStartOfDay(java.time.ZoneOffset.UTC).toInstant.toEpochMilli)
273273
case StringType => generator.writeString(value.toString)
274274
case _ => {
275275
val className = schema.getClass().getName()

spark/sql-35/src/test/scala/org/opensearch/spark/sql/DataFrameValueWriterTest.scala

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,10 +30,13 @@
3030
package org.opensearch.spark.sql
3131

3232
import java.io.ByteArrayOutputStream
33+
import java.sql.Date
34+
import java.util.TimeZone
3335

3436
import org.apache.spark.sql.Row
3537
import org.apache.spark.sql.catalyst.ScalaReflection
3638
import org.apache.spark.sql.types.ArrayType
39+
import org.apache.spark.sql.types.DateType
3740
import org.apache.spark.sql.types.IntegerType
3841
import org.apache.spark.sql.types.MapType
3942
import org.apache.spark.sql.types.StringType
@@ -165,4 +168,18 @@ class DataFrameValueWriterTest {
165168
}
166169
}
167170

171+
@Test
172+
def testWriteDateUsesUtcStartOfDay(): Unit = {
173+
val defaultTimeZone = TimeZone.getDefault
174+
try {
175+
TimeZone.setDefault(TimeZone.getTimeZone("Asia/Tokyo"))
176+
val schema = StructType(Seq(StructField("d", DateType)))
177+
val row = Row(Date.valueOf("2023-07-22"))
178+
val serialized = serialize(row, schema)
179+
assertTrue(serialized.contains(""""d":1689984000000"""))
180+
} finally {
181+
TimeZone.setDefault(defaultTimeZone)
182+
}
183+
}
184+
168185
}

spark/sql-40/src/main/scala/org/opensearch/spark/sql/DataFrameValueWriter.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -269,7 +269,7 @@ class DataFrameValueWriter(writeUnknownTypes: Boolean = false) extends Filtering
269269
case DoubleType => generator.writeNumber(value.asInstanceOf[Double])
270270
case FloatType => generator.writeNumber(value.asInstanceOf[Float])
271271
case TimestampType => generator.writeNumber(value.asInstanceOf[Timestamp].getTime())
272-
case DateType => generator.writeNumber(value.asInstanceOf[Date].getTime())
272+
case DateType => generator.writeNumber(value.asInstanceOf[Date].toLocalDate.atStartOfDay(java.time.ZoneOffset.UTC).toInstant.toEpochMilli)
273273
case StringType => generator.writeString(value.toString)
274274
case _ => {
275275
val className = schema.getClass().getName()

spark/sql-40/src/test/scala/org/opensearch/spark/sql/DataFrameValueWriterTest.scala

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,10 +30,13 @@
3030
package org.opensearch.spark.sql
3131

3232
import java.io.ByteArrayOutputStream
33+
import java.sql.Date
34+
import java.util.TimeZone
3335

3436
import org.apache.spark.sql.Row
3537
import org.apache.spark.sql.catalyst.ScalaReflection
3638
import org.apache.spark.sql.types.ArrayType
39+
import org.apache.spark.sql.types.DateType
3740
import org.apache.spark.sql.types.IntegerType
3841
import org.apache.spark.sql.types.MapType
3942
import org.apache.spark.sql.types.StringType
@@ -165,4 +168,18 @@ class DataFrameValueWriterTest {
165168
}
166169
}
167170

171+
@Test
172+
def testWriteDateUsesUtcStartOfDay(): Unit = {
173+
val defaultTimeZone = TimeZone.getDefault
174+
try {
175+
TimeZone.setDefault(TimeZone.getTimeZone("Asia/Tokyo"))
176+
val schema = StructType(Seq(StructField("d", DateType)))
177+
val row = Row(Date.valueOf("2023-07-22"))
178+
val serialized = serialize(row, schema)
179+
assertTrue(serialized.contains(""""d":1689984000000"""))
180+
} finally {
181+
TimeZone.setDefault(defaultTimeZone)
182+
}
183+
}
184+
168185
}

0 commit comments

Comments
 (0)