如何在Databricks中用SQL/SparkSQL/PySpark更新Delta表test_table
在Databricks中实现Delta表与CSV数据的匹配更新
问题背景
Databricks中存在Delta表 test_table,数据如下:
data_value,batch_number,rnumber,file_name,responseStatus,url 21gdsddnh,1,1,subject_1.csv,'','','' 21gdsdernh,1,2,subject_1.csv,'','','' 21gdweddnh,1,3,subject_1.csv,'','','' 21gdweddnh,2,1,subject_2.csv,'','','' 21gdweddnh,2,2,subject_2.csv,'','','' 21gdweddnh,2,3,subject_2.csv,'','',''
另有存储路径下的response.csv文件,内容如下:
file_name,responseStatus,url subject_1.csv,SUCCESS,/org/source/restapi.com subject_1.csv,SUCCESS,/org/source/restapi.com subject_1.csv,SUCCESS,/org/source/restapi.com subject_2.csv,SUCCESS,/org/source/restapi.com subject_2.csv,SUCCESS,/org/source/restapi.com subject_2.csv,SUCCESS,/org/source/restapi.com
需求:将response.csv的6行数据按对应关系匹配更新到test_table,填充responseStatus和url字段。
方法一:SparkSQL实现
操作步骤
- 加载
response.csv为临时视图,按file_name分区添加行号,用于和test_table的rnumber匹配 - 利用Delta Lake的
MERGE INTO语法完成匹配更新
代码示例
-- 加载response.csv为带行号的临时视图 CREATE OR REPLACE TEMP VIEW response_view AS SELECT file_name, responseStatus, url, ROW_NUMBER() OVER (PARTITION BY file_name ORDER BY (SELECT NULL)) AS rn FROM csv.`/your/path/to/response.csv` -- 替换为实际文件路径 -- 执行MERGE更新 MERGE INTO test_table t USING response_view r ON t.file_name = r.file_name AND t.rnumber = r.rn WHEN MATCHED THEN UPDATE SET t.responseStatus = r.responseStatus, t.url = r.url
方法二:PySpark实现
操作步骤
- 读取
response.csv并按file_name分区添加行号 - 关联
test_table与带行号的CSV数据,填充目标字段 - 覆盖原Delta表或使用Delta Merge操作完成更新
代码示例
from pyspark.sql.window import Window from pyspark.sql.functions import row_number # 1. 读取response.csv并添加匹配用的行号 response_df = spark.read.csv("/your/path/to/response.csv", header=True, inferSchema=True) window_spec = Window.partitionBy("file_name").orderBy("file_name") response_df = response_df.withColumn("rn", row_number().over(window_spec)) # 方式1:关联后覆盖原表 test_df = spark.read.table("test_table") updated_df = test_df.join( response_df, (test_df.file_name == response_df.file_name) & (test_df.rnumber == response_df.rn), "inner" ).select( test_df.data_value, test_df.batch_number, test_df.rnumber, test_df.file_name, response_df.responseStatus, response_df.url ) updated_df.write.mode("overwrite").format("delta").saveAsTable("test_table") # 方式2:使用Delta Merge API(推荐用于增量更新场景) from delta.tables import DeltaTable delta_table = DeltaTable.forName(spark, "test_table") delta_table.alias("t").merge( response_df.alias("r"), "t.file_name = r.file_name AND t.rnumber = r.rn" ).whenMatchedUpdate(set={ "responseStatus": "r.responseStatus", "url": "r.url" }).execute()
内容的提问来源于stack exchange,提问作者AzSurya Teja
相关产品推荐
相关产品推荐

