Skip to content

Commit

Permalink
[HUDI-2642] Add support ignoring case in update sql operation (#3882)
Browse files Browse the repository at this point in the history
  • Loading branch information
dongkelun committed Nov 30, 2021
1 parent 3433f00 commit a398aad
Show file tree
Hide file tree
Showing 2 changed files with 47 additions and 3 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,10 @@ case class UpdateHoodieTableCommand(updateTable: UpdateTable) extends RunnableCo
}.toMap

val updateExpressions = table.output
.map(attr => name2UpdateValue.getOrElse(attr.name, attr))
.map(attr => {
val UpdateValueOption = name2UpdateValue.find(f => sparkSession.sessionState.conf.resolver(f._1, attr.name))
if(UpdateValueOption.isEmpty) attr else UpdateValueOption.get._2
})
.filter { // filter the meta columns
case attr: AttributeReference =>
!HoodieRecord.HOODIE_META_COLUMNS.asScala.toSet.contains(attr.name)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,47 @@ package org.apache.spark.sql.hudi
class TestUpdateTable extends TestHoodieSqlBase {

test("Test Update Table") {
withTempDir { tmp =>
Seq("cow", "mor").foreach {tableType =>
val tableName = generateTableName
// create table
spark.sql(
s"""
|create table $tableName (
| ID int,
| NAME string,
| PRICE double,
| TS long
|) using hudi
| location '${tmp.getCanonicalPath}/$tableName'
| options (
| type = '$tableType',
| primaryKey = 'ID',
| preCombineField = 'TS'
| )
""".stripMargin)
// insert data to table
spark.sql(s"insert into $tableName select 1, 'a1', 10, 1000")
checkAnswer(s"select id, name, price, ts from $tableName")(
Seq(1, "a1", 10.0, 1000)
)

// update data
spark.sql(s"update $tableName set price = 20 where id = 1")
checkAnswer(s"select id, name, price, ts from $tableName")(
Seq(1, "a1", 20.0, 1000)
)

// update data
spark.sql(s"update $tableName set price = price * 2 where id = 1")
checkAnswer(s"select id, name, price, ts from $tableName")(
Seq(1, "a1", 40.0, 1000)
)
}
}
}

test("Test ignoring case for Update Table") {
withTempDir { tmp =>
Seq("cow", "mor").foreach {tableType =>
val tableName = generateTableName
Expand All @@ -46,13 +87,13 @@ class TestUpdateTable extends TestHoodieSqlBase {
)

// update data
spark.sql(s"update $tableName set price = 20 where id = 1")
spark.sql(s"update $tableName set PRICE = 20 where ID = 1")
checkAnswer(s"select id, name, price, ts from $tableName")(
Seq(1, "a1", 20.0, 1000)
)

// update data
spark.sql(s"update $tableName set price = price * 2 where id = 1")
spark.sql(s"update $tableName set PRICE = PRICE * 2 where ID = 1")
checkAnswer(s"select id, name, price, ts from $tableName")(
Seq(1, "a1", 40.0, 1000)
)
Expand Down

0 comments on commit a398aad

Please sign in to comment.