如何用PySpark XML连接器重命名输出文件并移除_SUCCESS文件
问题描述
使用PySpark XML连接器写入XML文件时,已成功生成分块文件(如part-00000),需要实现两个需求:
- 将分块XML文件重命名为当前时间戳命名的文件(例如
20240520143000.xml) - 删除输出目录中的
_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}")
关键细节说明
- 时间戳生成:用
datetime.now().strftime("%Y%m%d%H%M%S")生成精确到秒的时间戳,确保文件名唯一且易识别,可根据需求调整格式(如加入时区、毫秒)。 - 文件匹配逻辑:通过
startswith("part-")和endswith(".xml")定位目标分块文件,因为已经用coalesce(1)限制了输出为单个分块,所以无需处理多文件情况。 - 分布式文件系统兼容:如果操作HDFS等分布式存储,
os模块无法直接使用,需改用pyarrow.hdfs或HDFS客户端库完成文件操作,示例逻辑可平移适配。 - 性能提示:
coalesce(1)仅适合小数据量场景,大数据量下会导致单节点压力过大,需谨慎使用。
内容的提问来源于stack exchange,提问作者anuj
相关产品推荐
相关产品推荐

