apache / apache/hudi

If Sanitastiion Enabled In HudiStreamer It is taking too much time

Open
#10,466 6 comments 0 reactions 0 assignees View on GitHub
area:ingest priority:high
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

**_Tips before filing an issue_**

**Describe the problem you faced**

I have enabled the SANITIZE_SCHEMA_FIELD_NAMES hudiDeltaStreamer is stuck after reading CSV.
I think we can refactor the code it too better way.
Instead of using withColumnRenamed the transformation should be something like this

def transformSchemaBeginEndCharReplace(spark: SparkSession, final_stream: Dataset[Row], pii_masking_col: Seq[Any]): Dataset[Row] = {
val sql_select = new StringBuilder
val schema = final_stream.schema
for (i <- schema) {
if (i.dataType.isInstanceOf[StructType] || i.dataType.isInstanceOf[ArrayType]) {
sql_select.append(s"cast(to_json(`${i.name}`) as String)")
sql_select.append(" as ")
sql_select.append(avroSchemaNameConversionBeginEndCharReplace(i.name) + " , ")

}
else if (pii_masking_col.contains(i.name)) {
sql_select.append(s"sha1(`${i.name}`)")
sql_select.append(" as ")
sql_select.append(avroSchemaNameConversionBeginEndCharReplace(i.name) + " , ")
}
else {
sql_select.append(s"`${i.name}`")
sql_select.append(" as ")
sql_select.append(avroSchemaNameConversionBeginEndCharReplace(i.name) + " , ")

}
}
val final_sql = sql_select.toString().stripSuffix(" , ").split(",")
final_stream.selectExpr(final_sql: _*)
}

def avroSchemaNameConversionBeginEndCharReplace(name: String) = {
val regexPattern = "(^[0-9])|(^[^a-zA-Z_])|(([^A-Za-z0-9_])$)|([^A-Za-z0-9_])".r
val outputString = regexPattern.replaceAllIn(name, m => {
if(m.group(1)!=null){
s"_${m.group(1)}"
}
else if(m.group(2)!=null || m.group(3) != null ){
""
}
else {
"_"
}
})

outputString
}

We can set and adjust this work faster for my local transformation

**To Reproduce**

Steps to reproduce the behavior:

1.
2.
3.
4.

**Expected behavior**

A clear and concise description of what you expected to happen.

**Environment Description**

* Hudi version : 0.14.1

* Spark version : 3.3

* Hive version :

* Hadoop version :

* Storage (HDFS/S3/GCS..) : S3

* Running on Docker? (yes/no) :

**Additional context**

Add any other context about the problem here.

**Stacktrace**

```Add the stacktrace of the error.```

Contributor guide

No contributing guide indexed for this repository

Research direction

Start by reproducing the delay in HudiStreamer with SANITIZE_SCHEMA_FIELD_NAMES enabled, using the reported Hudi 0.14.1, Spark 3.3, CSV input, and S3 environment. Compare the existing withColumnRenamed transformation with the proposed selectExpr approach; done means CSV processing completes without the reported excessive delay while schema-name sanitization and masking remain correct.

Written by the indexing model from the issue text.

Assessment

Tech stack
scala
Domain
data-engineering, performance, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.