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

更新Parquet文件格式相关咨询:是否需用Spark DataFrames执行Upserts?

关于ADLS中Parquet数据迁移与Upserts的问题解答

是否需要使用Spark DataFrames?

是的,绝大多数场景下用Spark DataFrames是最优选择:

  • Spark对Parquet格式的读写支持成熟,能借助分布式处理能力高效应对ADLS上的大规模数据;
  • 不管是简单复制、字段过滤/重命名、聚合等操作,都需要将Parquet文件加载为DataFrame来处理;
  • 若仅做文件级复制(不修改数据内容),可使用ADLS原生工具(如azcopy)直接复制,但这类场景极少——迁移Parquet数据通常伴随数据处理需求。

是否需要执行Upserts操作?

这完全取决于你的具体需求:

  • 无需Upserts的场景:
    • 目标文件夹是全新的,无已有数据;
    • 你希望完全替换目标文件夹中的现有数据。
      这种情况直接读取源数据、处理后用mode="overwrite"写入目标路径即可,不用做Upserts。
  • 需要Upserts的场景:
    • 目标文件夹已有Parquet数据,且你需要合并新数据与原有数据(比如更新重复主键的记录、插入新记录)。
      注意:Parquet是列式存储,不支持原地更新,Upserts必须通过「读入新旧数据→合并处理→重新写入」的方式实现。

代码示例

常规读写(无Upserts需求)

from pyspark.sql import SparkSession

# 初始化Spark会话
spark = SparkSession.builder.appName("ParquetADLSMigration").getOrCreate()

# 读取ADLS源Parquet文件
source_df = spark.read.parquet("abfss://<container>@<storage-account>.dfs.core.windows.net/path/to/source")

# 数据处理示例:过滤字段、重命名
processed_df = source_df.select("user_id", "order_amount", "order_date").withColumnRenamed("order_amount", "amount")

# 写入目标ADLS路径,覆盖原有数据
processed_df.write.mode("overwrite").parquet("abfss://<container>@<storage-account>.dfs.core.windows.net/path/to/target")

Upserts合并数据(有新旧数据合并需求)

假设以user_id为主键,保留update_time最新的记录:

from pyspark.sql import SparkSession
from pyspark.sql.functions import when, col

spark = SparkSession.builder.appName("ParquetUpsert").getOrCreate()

# 读取源数据与目标现有数据
source_df = spark.read.parquet("abfss://<container>@<storage-account>.dfs.core.windows.net/path/to/new-data")
target_df = spark.read.parquet("abfss://<container>@<storage-account>.dfs.core.windows.net/path/to/existing-data")

# 合并数据:优先保留源数据中更新时间更晚的记录,同时保留目标中独有的记录
merged_df = source_df.join(target_df, on="user_id", how="full_outer") \
    .select(
        col("user_id"),
        when(col("source.update_time") > col("target.update_time"), col("source.amount")).otherwise(col("target.amount")).alias("amount"),
        when(col("source.update_time").isNotNull(), col("source.update_time")).otherwise(col("target.update_time")).alias("update_time")
    )

# 写入目标路径,覆盖原有数据
merged_df.write.mode("overwrite").parquet("abfss://<container>@<storage-account>.dfs.core.windows.net/path/to/target")

内容的提问来源于stack exchange,提问作者RData

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 08:01:34