【发布时间】:2022-07-28 18:43:47
【问题描述】:
我有一个增量表old,我想将它与new 合并。在new 表中,有一些id 值也存在于old 表中。我想通过总结old 和new 表cons 值来更新重叠ids 的cons 值。该怎么做?
【问题讨论】:
标签: pyspark delta-lake delta
我有一个增量表old,我想将它与new 合并。在new 表中,有一些id 值也存在于old 表中。我想通过总结old 和new 表cons 值来更新重叠ids 的cons 值。该怎么做?
【问题讨论】:
标签: pyspark delta-lake delta
试试这个:
在增量表更新中,您可以像在创建任何新的 spark 列时一样执行算术运算。
import pyspark.sql.functions as F
from delta.tables import *
spark.createDataFrame([{"id":i, "cons":1, "cons2":1} for i in range(500)])\
.write.format("delta").mode("overwrite").option("overwriteSchema", "true")\
.save("dbfs:/FileStore/anmol/sample_events_croma_before")
new = spark.createDataFrame([{"id":i, "cons":1, "cons2":1} for i in range(450, 550)])
old = DeltaTable.forPath(spark, "dbfs:/FileStore/anmol/sample_events_croma_before")
old.alias('old')\
.merge(new.alias('new')\
, "old.id = new.id")\
.whenMatchedUpdate(set={
"id": "new.id",
"cons": "old.cons + new.cons",
"cons2": F.col("old.cons2") + F.col("new.cons2"),
})\
.whenNotMatchedInsert(values={
"id": "new.id",
"cons": "new.cons",
})\
.execute()
【讨论】: