如何在PySpark中将Foundry数据集写入XML格式文件
Foundry平台写出XML格式数据集实现方案
你可以直接对接databricks:spark-xml的原生写接口,结合Foundry数据集的底层存储路径完成XML写入,支持单条记录对应单个XML文件的需求。
核心注意事项
- 确保Transform运行环境已预装
databricks:spark-xml连接器,版本和读取XML时使用的版本保持一致 - 禁止在写完XML文件后调用Output的
write_dataframe()方法,否则会覆盖已生成的XML文件 - 若需要每行数据对应一个独立XML文件,需要先对DataFrame按唯一标识重分区,保证单分区仅包含1条数据,写完后重命名默认生成的part分片文件即可
可直接运行的代码示例
from transforms.api import transform, Input, Output from pyspark.sql.types import StructField, StructType, StringType, DoubleType from pyspark.sql import functions as F import os BOOK_SCHEMA = StructType([ StructField("_id", StringType(), True), StructField("author", StringType(), True), StructField("description", StringType(), True), StructField("genre", StringType(), True), StructField("price", DoubleType(), True), StructField("publish_date", StringType(), True), StructField("title", StringType(), True)] ) @transform( source_df=Input("/output/book-xml"), xml_output=Output("/output/book-xml-files"), ) def compute(ctx, source_df, xml_output): spark = ctx.spark_session # 读取已解析完成的结构化书籍数据 df = source_df.dataframe("selected") # 以下为单条记录对应单个XML文件的可选逻辑,不需要可删除 # 为每条记录生成唯一XML文件名,按文件名重分区保证单分区单条数据 df = df.withColumn("file_name", F.concat(F.col("_id"), F.lit(".xml"))) df = df.repartition(1, "file_name") # 获取输出数据集的底层Hadoop存储路径 output_fs = xml_output.filesystem() output_path = output_fs.hadoop_path # 调用spark-xml原生写接口,rowTag参数需和读取XML时的配置保持一致 df.write.format('xml') \ .option("rowTag", "book") \ .option("rootTag", "catalog") \ .mode("overwrite") \ .save(output_path) # 以下为单文件重命名的可选逻辑,不需要可删除 hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration() Path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path fs = Path(output_path).getFileSystem(hadoop_conf) for file_status in fs.listStatus(Path(output_path)): file_path = file_status.getPath() file_name = file_path.getName() # 仅处理XML分片文件,跳过_SUCCESS标记文件、隐藏文件 if file_name.startswith("part-") and file_name.endswith(".xml"): part_df = spark.read.format("xml").option("rowTag", "book").load(file_path.toString()) target_name = part_df.select("file_name").first()[0] fs.rename(file_path, Path(os.path.join(output_path, target_name)))
补充说明
- 如果不需要单条记录对应单个XML文件,直接跳过重分区、重命名步骤即可,写出的分片XML文件可以直接被已有的XML读取逻辑正常识别
- XML写出支持自定义编码、XML版本声明、属性前缀等配置,直接在
.option()中传入对应spark-xml参数即可 - 写入完成后Foundry会自动识别数据集的文件更新,不需要额外调用提交方法
内容的提问来源于stack exchange,提问作者huy
相关产品推荐
相关产品推荐

