【问题标题】:Spark: How to write bytes string to hdfs hadoop in pyspark for spark-xml transformation?Spark:如何在 pyspark 中将字节字符串写入 hdfs hadoop 以进行 spark-xml 转换?
【发布时间】:2021-04-20 00:25:09
【问题描述】:

在python中,字节字符串可以简单地保存到单个xml文件中:

with open('/home/user/file.xml' ,'wb') as f:
    f.write(b'<Value>1</Value>') 
   

当前输出:/home/user/file.xml(保存在本地文件中的文件)

问题:如何在pyspark中将字符串保存到hdfs上的xml文件:

预期输出:'hdfs://hostname:9000/file.xml'

背景:大量 xml 文件由 3rd 方 Web API 提供。我在 pyspark 中构建到 delta 湖的 ETL 管道。数据由aiohttp异步提取,接下来我想在将spark数据帧保存到delta湖之前使用spark-xml进行转换(需要pyspark)。我正在寻找最有效的方式来构建管道。

在 github 上向 spark-xml 开发人员提出了类似的问题。 https://github.com/databricks/spark-xml/issues/515

最新研究:

  1. spark-xml 用作输入 xml 文件直接存储为磁盘上的文本或 spark 数据帧

  2. 所以我只能使用以下 2 个选项之一:

a) 一些 hdfs 客户端(pyarrow,hdfs,aiohdfs) 将文件保存到 hdfs (hdfs 上的文本文件不是很有效的格式)

b) 将数据加载到 spark-xml 转换的 spark 数据帧(delta Lake 的本机格式)

如果您有其他想法,请告诉我。

【问题讨论】:

  • 这些似乎是无关的问题。您在阅读文件时遇到什么问题?
  • 我已将示例修改为最小代码。我的问题是我不知道如何将字节串直接保存到单个 hdfs 文件中。我从 web api 将 xml 文件下载到 python 中的字节字符串。我可以将 file.xml 保存到 python 中的本地文件系统。我找不到如何使用 pyspark 将字符串 b'1' 直接保存到 hdfs 的方法,所以目前我先将文件保存到本地文件系统,然后将其从本地 fs 复制到 hdfs这是浪费资源。
  • 你不需要 Spark。您可以使用 WebHDFS 或其他 Python 库来编写文件。 stackoverflow.com/questions/47926758/python-write-to-hdfs-file 不过,你还是需要某种方式的本地文件...
  • 感谢您的快速回答,我将测试建议的库,我有大型 ETL,我正在将文件从不同格式的不同 api 下载到 pyspark 中的 databricks delta 湖,我正在尝试保留尽可能少的库,所以当所有代码都在 pyspark 中时,我尝试直接在 pyspark 中对其进行编码。我使用 asyncio 进行下载,所以我也在尝试 aiohdfs。

标签: python hadoop hdfs


【解决方案1】:

不要被 databricks spark-xml 文档误导,因为它们会导致使用未压缩的 XML 文件作为输入。这是非常低效的,而且更快的是直接下载 XML 到 spark 数据帧。 Databricks xml-pyspark 版本不包含但有一个workaround:

from pyspark.sql.column import Column, _to_java_column
from pyspark.sql.types import _parse_datatype_json_string

def ext_from_xml(xml_column, schema, options={}):
    java_column = _to_java_column(xml_column.cast('string'))
    java_schema = spark._jsparkSession.parseDataType(schema.json())
    scala_map = spark._jvm.org.apache.spark.api.python.PythonUtils.toScalaMap(options)
    jc = spark._jvm.com.databricks.spark.xml.functions.from_xml(
        java_column, java_schema, scala_map)
    return Column(jc)

def ext_schema_of_xml_df(df, options={}):
    assert len(df.columns) == 1

    scala_options = spark._jvm.PythonUtils.toScalaMap(options)
    java_xml_module = getattr(getattr(
        spark._jvm.com.databricks.spark.xml, "package$"), "MODULE$")
    java_schema = java_xml_module.schema_of_xml_df(df._jdf, scala_options)
    return _parse_datatype_json_string(java_schema.json())

已下载到列表的 XMLs

xml = [('url',"""<Level_0 Id0="Id0_value_file1">
    <Level_1 Id1_1 ="Id3_value" Id_2="Id2_value">
      <Level_2_A>A</Level_2_A>
      <Level_2>
        <Level_3>
          <Level_4>
            <Date>2021-01-01</Date>
            <Value>4_1</Value>
          </Level_4>
          <Level_4>
            <Date>2021-01-02</Date>
            <Value>4_2</Value>
          </Level_4>
        </Level_3>
      </Level_2>
    </Level_1>
  </Level_0>"""),

  ('url',"""<Level_0 I"d0="Id0_value_file2">
    <Level_1 Id1_1 ="Id3_value" Id_2="Id2_value">
      <Level_2_A>A</Level_2_A>
      <Level_2>
        <Level_3>
          <Level_4>
            <Date>2021-01-01</Date>
            <Value>4_1</Value>
          </Level_4>
          <Level_4>
            <Date>2021-01-02</Date>
            <Value>4_2</Value>
          </Level_4>
        </Level_3>
      </Level_2>
    </Level_1>
  </Level_0>""")]

XML字符串的Spark数据框转换:

#create df with XML strings  
 rdd = sc.parallelize(xml)
 df = spark.createDataFrame(rdd,"url string, content string")

# XML schema
 payloadSchema = ext_schema_of_xml_df(df.select("content"))

 # parse xml
 parsed = df.withColumn("parsed", ext_from_xml(df.content, payloadSchema, {"rowTag":"Level_0"}))

# select required data
  df2 = parsed.select(
    'parsed._Id0',
    F.explode_outer('parsed.Level_1.Level_2.Level_3.Level_4').alias('Level_4')
  ).select(
      '`parsed._Id0`',
      'Level_4.*'
  )

解码字节:b'string'.decode('utf-8')

@mck 回答有关 XML 的更多信息: How to transform to spark Data Frame data from multiple nested XML files with attributes

【讨论】:

    猜你喜欢
    • 2018-07-29
    • 2020-08-03
    • 2018-11-16
    • 1970-01-01
    • 2015-11-27
    • 1970-01-01
    • 2022-01-17
    • 2021-12-25
    • 2016-08-03
    相关资源
    最近更新 更多