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


##########
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

Review Comment:
   Addressed in 4eab0eb90767. Repointed DATE_STRING at separate, non-default 
input/output timezones (GMT+08:00 to GMT-05:00), including a date rollover, so 
it adds coverage beyond the existing same-zone test.



##########
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)

Review Comment:
   Addressed in 4eab0eb90767. Replaced hours with lowercase microseconds using 
a fixed microsecond timestamp. A separate Turkish-locale regression test 
exercises the locale-sensitive normalization.



##########
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) {

Review Comment:
   Addressed in 4eab0eb90767. Kept only the timezone-split DATE_STRING and 
lowercase-microseconds SCALAR cases: four write/read cycles across AVRO and 
SPARK instead of eight. Each case now writes two rows to one deterministic 
partition. Both AVRO and SPARK integration-test invocations passed locally.



-- 
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