brightwon opened a new issue, #10995:
URL: https://github.com/apache/hudi/issues/10995
**Describe the problem you faced**
I'm operating a typical Hudi workload that involves using spark structured
streaming to read CDC events from Kafka and perform Upserts into S3.
I've encountered an issue where, when rows with same precombine values are
present in the same batch, applying repartition to a Kafka input dataframe
results saving rows that are not the last.
I faced a slow performance issue during the Tagging stage. To address this,
I applied repartition to the input dataframe, which indeed improved the
performance. Unfortunately, this led to a situation where data with incorrect
order was saved.
**To Reproduce**
Steps to reproduce the behavior:
1. Sent test data to Kafka using a Kafka producer. Here's an example of the
data:
Offset 1 :
```
{
"pk": 1,
"c1": 1,
"c2": 1
}
```
Offset 2 :
```
{
"pk": 1,
"c1": 1,
"c2": 2
}
```
produced 1000 records, incrementing the value of the c2 field by 1 for each
(1 ~ 1000). The offset key matches the value of the "pk" field in value, where
c1 is the precombine field, and c2 is the column used to distinguish each row.
2. Run the example spark application to save the data to S3.
code sample
```
val cdcDF = ss
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker1,broker2,broker3")
.option("topic", "hudi_inorder_test")
.option("maxOffsetsPerTrigger", "1000000")
.load()
val stream = cdcDF
.repartition(100) // apply repartition for tagging stage parallelism
.writeStream
.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
saveBatch(batchDF, batchId) }
.trigger(Trigger.ProcessingTime(60, TimeUnit.SECONDS))
.start()
stream.awaitTermination()
def saveBatch(batchDF: DataFrame, batchId: Long): Unit = {
batchDF.persist()
// ...deserialize the avro offset value...
val valueDF = batchDF.select(...)
// save hudi dataset
valueDF
.write
.format("hudi")
.options(hudiOptions)
.mode("append")
.save("s3://mybucket/test/")
batchDF.unpersist()
}
```
write config (hudiOptions)
```
{
"hoodie.table.name": "hudi_inorder_test",
"hoodie.datasource.write.table.type": "MERGE_ON_READ",
"hoodie.datasource.write.recordkey.field": "pk",
"hoodie.datasource.write.partitionpath.field": "",
"hoodie.datasource.write.precombine.field": "c1",
"hoodie.datasource.write.keygenerator.class":
"org.apache.hudi.keygen.CustomKeyGenerator",
"hoodie.compact.inline": "true",
"hoodie.datasource.compaction.async.enable": "false",
"hoodie.compact.inline.max.delta.commits": "1"
}
```
3. Check the saved results. (I used AWS Athena.)
4. Repeat steps 1 and 3 to verify if the value of c2 continues to change.
The c2 value changes with each repetition. here is the saved results.
1.
```
{
"_hoodie_commit_time": "20240405065000817",
"_hoodie_commit_seqno": "20240405065000817_0_1",
"_hoodie_record_key": "1",
"_hoodie_partition_path": "",
"_hoodie_file_name":
"1cee1e0b-2d25-434a-adb7-3d7e0ae4a275-0_0-58-9709_20240405065000817.parquet",
"pk": "1",
"c1": "1",
"c2": "496"
}
```
2.
```
{
"_hoodie_commit_time": "20240405065319747",
"_hoodie_commit_seqno": "20240405065319747_0_1",
"_hoodie_record_key": "1",
"_hoodie_partition_path": "",
"_hoodie_file_name":
"1cee1e0b-2d25-434a-adb7-3d7e0ae4a275-0_0-115-15330_20240405065319747.parquet",
"pk": "1",
"c1": "1",
"c2": "316"
}
```
If repartition() is removed from the code, the c2 value is correctly saved
as the last offset value, 1000.
Since I'll be using a single Kafka partition, I considered using the Kafka
offset number as the precombine field. However, Hudi does not support changes
to the precombine field.
I want to improve the performance of the tagging stage through increased
parallelism without causing issues with data ordering. How can I resolve this?
**Environment Description**
* Hudi version : 0.10.1-amzn-0 (EMR 6.6.0)
* Spark version : 3.2.0 (EMR 6.6.0)
* Hive version : 3.1.2 (EMR 6.6.0)
* Hadoop version : 3.2.1 (EMR 6.6.0)
* Storage (HDFS/S3/GCS..) : S3
* Running on Docker? (yes/no) : no
--
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]