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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 23:21:31