You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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实现

操作步骤

  1. 加载response.csv为临时视图,按file_name分区添加行号,用于和test_table的rnumber匹配
  2. 利用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实现

操作步骤

  1. 读取response.csv并按file_name分区添加行号
  2. 关联test_table与带行号的CSV数据,填充目标字段
  3. 覆盖原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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 14:40:15