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

如何用PySpark XML连接器重命名输出文件并移除_SUCCESS文件

问题描述

使用PySpark XML连接器写入XML文件时,已成功生成分块文件(如part-00000),需要实现两个需求:

  1. 将分块XML文件重命名为当前时间戳命名的文件(例如20240520143000.xml)
  2. 删除输出目录中的_SUCCESS文件

现有实现代码如下:

import os
from datetime import datetime
from pyspark.sql import SparkSession
from pyspark.sql.functions import struct

data = [
    (
        "xxx", "VN", "Goods", "bIKE", "Ignition", "20100", "13000",
        "1.0", "IT", "2024-01-23T13:15:30.45+01:00"
    ),
]

schema = [
    "PartNumber", "PartName", "Transport", "Vehicle", "Engine", "MaxWeight", "CylinderCapacity",
    "MsgVersion", "SenderID", "SendTime"
]

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

df_nested = spark.createDataFrame(data, schema=schema) \
    .withColumn("GmfHeader", struct("MsgVersion", "SenderID", "SendTime")) \
    .withColumn("Product", struct("Transport", "Vehicle", "Engine", "MaxWeight", "CylinderCapacity"))

df_nested = df_nested.drop("MsgVersion", "SenderID", "SendTime", "Transport", "Vehicle", "Engine", "MaxWeight", "CylinderCapacity")

output_path = "/mnt/test102/xml_output/"

df_nested.coalesce(1).write \
    .format("xml") \
    .option("rootTag", "n1:Part") \
    .option("rowTag", "n1:PartMastInf") \
    .mode("overwrite") \
    .save(output_path)

print(f"XML files generated successfully at: {output_path}")
解决方案

核心思路是在Spark写入完成后,遍历输出目录定位目标文件,完成重命名和_SUCCESS文件删除操作,具体实现如下:

修改后的完整代码

import os
from datetime import datetime
from pyspark.sql import SparkSession
from pyspark.sql.functions import struct

data = [
    (
        "xxx", "VN", "Goods", "bIKE", "Ignition", "20100", "13000",
        "1.0", "IT", "2024-01-23T13:15:30.45+01:00"
    ),
]

schema = [
    "PartNumber", "PartName", "Transport", "Vehicle", "Engine", "MaxWeight", "CylinderCapacity",
    "MsgVersion", "SenderID", "SendTime"
]

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

df_nested = spark.createDataFrame(data, schema=schema) \
    .withColumn("GmfHeader", struct("MsgVersion", "SenderID", "SendTime")) \
    .withColumn("Product", struct("Transport", "Vehicle", "Engine", "MaxWeight", "CylinderCapacity"))

df_nested = df_nested.drop("MsgVersion", "SenderID", "SendTime", "Transport", "Vehicle", "Engine", "MaxWeight", "CylinderCapacity")

output_path = "/mnt/test102/xml_output/"

# 写入XML文件
df_nested.coalesce(1).write \
    .format("xml") \
    .option("rootTag", "n1:Part") \
    .option("rowTag", "n1:PartMastInf") \
    .mode("overwrite") \
    .save(output_path)

# 生成带时间戳的目标文件名
timestamp = datetime.now().strftime("%Y%m%d%H%M%S")
target_filename = f"{timestamp}.xml"
target_file_path = os.path.join(output_path, target_filename)

# 遍历输出目录处理文件
for file_name in os.listdir(output_path):
    full_file_path = os.path.join(output_path, file_name)
    # 重命名part开头的XML文件
    if file_name.startswith("part-") and file_name.endswith(".xml"):
        os.rename(full_file_path, target_file_path)
        print(f"文件已重命名:{file_name} → {target_filename}")
    # 删除_SUCCESS文件
    elif file_name == "_SUCCESS":
        os.remove(full_file_path)
        print("已删除_SUCCESS标记文件")

print(f"最终XML文件路径:{target_file_path}")

关键细节说明

  1. 时间戳生成:用datetime.now().strftime("%Y%m%d%H%M%S")生成精确到秒的时间戳,确保文件名唯一且易识别,可根据需求调整格式(如加入时区、毫秒)。
  2. 文件匹配逻辑:通过startswith("part-")和endswith(".xml")定位目标分块文件,因为已经用coalesce(1)限制了输出为单个分块,所以无需处理多文件情况。
  3. 分布式文件系统兼容:如果操作HDFS等分布式存储,os模块无法直接使用,需改用pyarrow.hdfs或HDFS客户端库完成文件操作,示例逻辑可平移适配。
  4. 性能提示:coalesce(1)仅适合小数据量场景,大数据量下会导致单节点压力过大,需谨慎使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 15:45:38