如何在Delta Table的updateExpr中实现复杂字段更新逻辑
Delta Table Merge实现数组字段增量拼接
场景说明
- 增量更新Delta表时,
ts、user两个字段执行普通覆盖更新 scores字段为元素是Map的数组类型,需要将增量数据中的scores值与表内已有历史scores值拼接,不做直接覆盖替换
测试数据构造
val historicalDF = Seq( (1, 0, "Roger", Seq(Map("score" -> 5, "year" -> 2012))) ).toDF("id", "ts", "user", "scores") historicalDF.write .format("delta") .mode("overwrite") .save(table_path) val hist_dt : DeltaTable = DeltaTable.forPath(spark, table_path) val incrementalDF = Seq( (1, 1, "Roger Rabbit", Seq(Map("score" -> 7, "year" -> 2013))) ).toDF("id", "ts", "user", "scores")
预期输出
执行merge后预期得到如下结果(注:原示例预期结果存在笔误,拼接后应同时包含2012、2013年的分数记录):
+---+---+------------+--------------------------------------------------------+ |id |ts |user |scores | +---+---+------------+--------------------------------------------------------+ |1 |1 |Roger Rabbit|[{score -> 5, year -> 2012}, {score -> 7, year -> 2013}]| +---+---+------------+--------------------------------------------------------+
报错原因
原merge代码的核心问题是:updateExpr方法要求传入的Map值是Spark SQL表达式字符串,原写法"in.scores" ++ "ex.scores"是在Scala代码层面对两个字符串对象执行++拼接操作,逻辑根本没有传递到Spark执行层做数组运算,因此抛出value ++ is not a member of (String, String)错误。
实现方案
updateExpr支持所有合法的Spark SQL表达式,针对数组拼接场景,直接在表达式字符串中调用Spark SQL内置的concat函数即可完成合并,修正后的merge代码如下:
hist_dt .as("ex") .merge(incrementalDF.as("in"), "ex.id = in.id") .whenMatched .updateExpr( Map( "ts" -> "in.ts", "user" -> "in.user", "scores" -> "concat(in.scores, ex.scores)" // 调用SQL内置concat函数按顺序合并两个数组 ) ) .whenNotMatched .insertAll() .execute()
扩展说明:
- 如果需要对合并后的数组去重,可以将
concat替换为array_union函数,该函数会自动丢弃重复的数组元素- 其他更复杂的数组处理逻辑(比如过滤无效元素、结构转换等),都可以直接在表达式字符串中使用对应Spark SQL内置函数实现,无需额外编写UDF
内容的提问来源于stack exchange,提问作者Ignacio Alorre
相关产品推荐
相关产品推荐

