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

在ADLS Gen2中用Databricks PySpark原地修改CSV日期值

PySpark代码:在Databricks中处理ADLS Gen2的CSV文件并模拟原地修改日期字段

关键说明

ADLS Gen2属于对象存储,不支持文件的原生原地修改操作。我们可以通过读取原文件→处理数据→删除原文件→写入处理后的数据到原路径的流程,实现类似“原地修改”的效果。

完整代码实现

# 1. 配置ADLS Gen2访问(若已通过Databricks工作区配置好权限可跳过此段)
# 替换为你的存储账户和Service Principal信息
spark.conf.set("fs.azure.account.auth.type.<your-storage-account>.dfs.core.windows.net", "OAuth")
spark.conf.set("fs.azure.account.oauth.provider.type.<your-storage-account>.dfs.core.windows.net", "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider")
spark.conf.set("fs.azure.account.oauth2.client.id.<your-storage-account>.dfs.core.windows.net", "<client-id>")
spark.conf.set("fs.azure.account.oauth2.client.secret.<your-storage-account>.dfs.core.windows.net", "<client-secret>")
spark.conf.set("fs.azure.account.oauth2.client.endpoint.<your-storage-account>.dfs.core.windows.net", "https://login.microsoftonline.com/<tenant-id>/oauth2/token")

# 2. 定义路径
source_path = "abfss://<container-name>@<your-storage-account>.dfs.core.windows.net/<csv-folder-path>/"

# 3. 读取CSV文件(替换`datetime_col`为你的实际日期字段名)
df = spark.read.csv(
    source_path,
    header=True,  # 若CSV无表头则设为False
    inferSchema=False  # 保持字符串类型便于处理
)

# 4. 提取日期部分,三种方法选其一即可
# 方法1:按空格分割取第一部分(适用于日期时间用空格分隔的场景)
df_processed = df.withColumn("datetime_col", df["datetime_col"].split(" ")[0])

# 方法2:截取前8位字符(适用于日期固定为MM-dd-yy格式的场景)
# df_processed = df.withColumn("datetime_col", df["datetime_col"].substr(1, 8))

# 方法3:用日期函数解析后格式化(更严谨,可处理格式异常)
# from pyspark.sql.functions import to_date, date_format
# df_processed = df.withColumn("datetime_col", date_format(to_date(df["datetime_col"], "MM-dd-yy h:mm:ss a"), "MM-dd-yy"))

# 5. 删除原路径下的CSV文件(跳过系统生成的_$folder$文件)
fs = spark.sparkContext._jvm.org.apache.hadoop.fs.FileSystem.get(spark.sparkContext._jsc.hadoopConfiguration())
path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(source_path)
for file_status in fs.listStatus(path):
    file_path = file_status.getPath().toString()
    if file_path.endswith(".csv") and not file_path.endswith("_$folder$"):
        fs.delete(spark.sparkContext._jvm.org.apache.hadoop.fs.Path(file_path), True)

# 6. 将处理后的数据写入原路径
df_processed.write.csv(
    source_path,
    header=True,
    mode="overwrite",
    quote='"',
    escape='"'
)

注意事项

  • 替换代码中所有<>包裹的占位符为你的实际信息;
  • 优先选择适配你数据格式的日期处理方法,方法3虽严谨但性能略低于前两种;
  • 执行前建议在测试环境验证,避免误删生产数据;
  • 确保Databricks集群拥有ADLS Gen2对应路径的读写权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 19:56:08