voonhous commented on code in PR #19657:
URL: https://github.com/apache/hudi/pull/19657#discussion_r3941462549


##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala:
##########
@@ -1450,6 +1451,74 @@ class TestCOWDataSource extends 
HoodieSparkClientTestBase with ScalaAssertionSup
     assertTrue(recordsReadDF.filter(col("_hoodie_partition_path") =!= 
udf_date_format(col("current_ts"))).count() == 0)
   }
 
+  @ParameterizedTest
+  @EnumSource(value = classOf[HoodieRecordType], names = Array("AVRO", 
"SPARK"))
+  def testTimestampBasedKeyGeneratorWithVariousConfigurations(recordType: 
HoodieRecordType) {
+    val (writeOpts, readOpts) = 
getWriterReaderOptsLessPartitionPath(recordType)
+
+    val records = recordsToStrings(dataGen.generateInserts("000", 
100)).asScala.toList
+    val inputDF = spark.read.json(spark.sparkContext.parallelize(records, 2))
+      .withColumn("current_ts_micros", col("current_ts") * 1000)
+      .withColumn("current_date_string",
+        date_format((col("current_ts") / 1000).cast("timestamp"), "yyyy-MM-dd 
HH:mm:ss"))
+      .withColumn("current_ts_hours", (col("current_ts") / 
3600000).cast("long"))
+
+    case class TestCase(partitionCol: String, tsType: String, outFmt: String,
+                        extraOpts: Map[String, String] = Map.empty,
+                        expectedPartitionUdf: 
org.apache.spark.sql.expressions.UserDefinedFunction)
+
+    def runTestCase(tc: TestCase): Unit = {
+      val writer = tc.extraOpts.foldLeft(
+        inputDF.write.format("hudi")
+          .options(writeOpts)
+          .option(KEYGENERATOR_CLASS_NAME.key(), 
classOf[TimestampBasedKeyGenerator].getName)
+          .mode(SaveMode.Overwrite)
+      ) { case (w, (k, v)) => w.option(k, v) }
+      writer.partitionBy(tc.partitionCol)
+        .option(TIMESTAMP_TYPE_FIELD.key, tc.tsType)
+        .option(TIMESTAMP_OUTPUT_DATE_FORMAT.key, tc.outFmt)
+        .save(basePath)
+      val readDF = 
spark.read.format("org.apache.hudi").options(readOpts).load(basePath)
+      assertTrue(readDF.filter(col("_hoodie_partition_path") =!= 
tc.expectedPartitionUdf(col(tc.partitionCol))).count() == 0)
+    }
+
+    val outputDateFmt = "yyyy-MM-dd HH"
+    // Joda's DateTimeZone.forID does not recognise "GMT+08:00". 
HoodieDateTimeParser resolves the
+    // configured id via java.util.TimeZone, so the expected values are built 
the same way.
+    val tzId = "GMT+08:00"
+
+    // Test 1: EPOCHMILLISECONDS with timezone GMT+08:00
+    val udfMillisTz = udf((millis: Long) => {
+      val zone = DateTimeZone.forTimeZone(TimeZone.getTimeZone(tzId))
+      new DateTime(millis, 
zone).toString(DateTimeFormat.forPattern(outputDateFmt).withZone(zone))
+    })
+    runTestCase(TestCase("current_ts", "EPOCHMILLISECONDS", outputDateFmt,
+      Map(TIMESTAMP_TIMEZONE_FORMAT.key -> tzId), udfMillisTz))
+
+    // Test 2: EPOCHMICROSECONDS (no timezone configured, so the key generator 
uses the JVM default)
+    val udfMicros = udf((micros: Long) =>
+      new DateTime(micros / 
1000).toString(DateTimeFormat.forPattern(outputDateFmt)))
+    runTestCase(TestCase("current_ts_micros", "EPOCHMICROSECONDS", 
outputDateFmt,
+      expectedPartitionUdf = udfMicros))
+
+    // Test 3: DATE_STRING with timezone
+    val dateStrInFmt = "yyyy-MM-dd HH:mm:ss"
+    val udfDateStrTz = udf((s: String) => {
+      val zone = DateTimeZone.forTimeZone(TimeZone.getTimeZone(tzId))
+      DateTime.parse(s, DateTimeFormat.forPattern(dateStrInFmt).withZone(zone))
+        .toString(DateTimeFormat.forPattern(outputDateFmt).withZone(zone))
+    })
+    runTestCase(TestCase("current_date_string", "DATE_STRING", outputDateFmt,
+      Map(TIMESTAMP_INPUT_DATE_FORMAT.key -> dateStrInFmt,
+        TIMESTAMP_TIMEZONE_FORMAT.key -> tzId), udfDateStrTz))

Review Comment:
   Verified at 4eab0eb9. Input `GMT+08:00` / output `GMT-05:00` now shift both 
the date and the hour, so the assertion fails if either zone is dropped. I ran 
the key generator on `2009-02-14 07:31:30` under `UTC`, `Asia/Singapore` and 
`America/New_York` and it returns `2009-02-13 18` in all three, so the expected 
value is JVM-zone independent too. Using the non-deprecated config pair also 
covers ground `testPartitionColumnsProperHandling` does not.



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala:
##########
@@ -1450,6 +1451,74 @@ class TestCOWDataSource extends 
HoodieSparkClientTestBase with ScalaAssertionSup
     assertTrue(recordsReadDF.filter(col("_hoodie_partition_path") =!= 
udf_date_format(col("current_ts"))).count() == 0)
   }
 
+  @ParameterizedTest
+  @EnumSource(value = classOf[HoodieRecordType], names = Array("AVRO", 
"SPARK"))
+  def testTimestampBasedKeyGeneratorWithVariousConfigurations(recordType: 
HoodieRecordType) {
+    val (writeOpts, readOpts) = 
getWriterReaderOptsLessPartitionPath(recordType)
+
+    val records = recordsToStrings(dataGen.generateInserts("000", 
100)).asScala.toList
+    val inputDF = spark.read.json(spark.sparkContext.parallelize(records, 2))
+      .withColumn("current_ts_micros", col("current_ts") * 1000)
+      .withColumn("current_date_string",
+        date_format((col("current_ts") / 1000).cast("timestamp"), "yyyy-MM-dd 
HH:mm:ss"))
+      .withColumn("current_ts_hours", (col("current_ts") / 
3600000).cast("long"))
+
+    case class TestCase(partitionCol: String, tsType: String, outFmt: String,
+                        extraOpts: Map[String, String] = Map.empty,
+                        expectedPartitionUdf: 
org.apache.spark.sql.expressions.UserDefinedFunction)
+
+    def runTestCase(tc: TestCase): Unit = {
+      val writer = tc.extraOpts.foldLeft(
+        inputDF.write.format("hudi")
+          .options(writeOpts)
+          .option(KEYGENERATOR_CLASS_NAME.key(), 
classOf[TimestampBasedKeyGenerator].getName)
+          .mode(SaveMode.Overwrite)
+      ) { case (w, (k, v)) => w.option(k, v) }
+      writer.partitionBy(tc.partitionCol)
+        .option(TIMESTAMP_TYPE_FIELD.key, tc.tsType)
+        .option(TIMESTAMP_OUTPUT_DATE_FORMAT.key, tc.outFmt)
+        .save(basePath)
+      val readDF = 
spark.read.format("org.apache.hudi").options(readOpts).load(basePath)
+      assertTrue(readDF.filter(col("_hoodie_partition_path") =!= 
tc.expectedPartitionUdf(col(tc.partitionCol))).count() == 0)
+    }
+
+    val outputDateFmt = "yyyy-MM-dd HH"
+    // Joda's DateTimeZone.forID does not recognise "GMT+08:00". 
HoodieDateTimeParser resolves the
+    // configured id via java.util.TimeZone, so the expected values are built 
the same way.
+    val tzId = "GMT+08:00"
+
+    // Test 1: EPOCHMILLISECONDS with timezone GMT+08:00
+    val udfMillisTz = udf((millis: Long) => {
+      val zone = DateTimeZone.forTimeZone(TimeZone.getTimeZone(tzId))
+      new DateTime(millis, 
zone).toString(DateTimeFormat.forPattern(outputDateFmt).withZone(zone))
+    })
+    runTestCase(TestCase("current_ts", "EPOCHMILLISECONDS", outputDateFmt,
+      Map(TIMESTAMP_TIMEZONE_FORMAT.key -> tzId), udfMillisTz))
+
+    // Test 2: EPOCHMICROSECONDS (no timezone configured, so the key generator 
uses the JVM default)
+    val udfMicros = udf((micros: Long) =>
+      new DateTime(micros / 
1000).toString(DateTimeFormat.forPattern(outputDateFmt)))
+    runTestCase(TestCase("current_ts_micros", "EPOCHMICROSECONDS", 
outputDateFmt,
+      expectedPartitionUdf = udfMicros))
+
+    // Test 3: DATE_STRING with timezone
+    val dateStrInFmt = "yyyy-MM-dd HH:mm:ss"
+    val udfDateStrTz = udf((s: String) => {
+      val zone = DateTimeZone.forTimeZone(TimeZone.getTimeZone(tzId))
+      DateTime.parse(s, DateTimeFormat.forPattern(dateStrInFmt).withZone(zone))
+        .toString(DateTimeFormat.forPattern(outputDateFmt).withZone(zone))
+    })
+    runTestCase(TestCase("current_date_string", "DATE_STRING", outputDateFmt,
+      Map(TIMESTAMP_INPUT_DATE_FORMAT.key -> dateStrInFmt,
+        TIMESTAMP_TIMEZONE_FORMAT.key -> tzId), udfDateStrTz))
+
+    // Test 4: SCALAR with hours (no timezone configured, so the key generator 
uses the JVM default)
+    val udfScalarHours = udf((hours: Long) =>
+      new 
DateTime(TimeUnit.HOURS.toMillis(hours)).toString(DateTimeFormat.forPattern(outputDateFmt)))
+    runTestCase(TestCase("current_ts_hours", "SCALAR", outputDateFmt,
+      Map(INPUT_TIME_UNIT.key -> "hours"), udfScalarHours))

Review Comment:
   Verified at 4eab0eb9. I compiled the patched 
`TimestampBasedAvroKeyGenerator` and confirmed 
`String.toUpperCase(Locale.ROOT)` in the bytecode; with it on the classpath, 
`microseconds` resolves under `-Duser.language=tr -Duser.country=TR`, where the 
previous build threw `No enum constant 
java.util.concurrent.TimeUnit.M<dotted-I>CROSECONDS`. 
`testScalarMicrosecondsWithTurkishLocale` pins it across GenericRecord, Row and 
InternalRow and restores the default locale in a `finally`, so it fails without 
the one-line fix.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to