【问题标题】:Dataflow - XML Source - Python - How?数据流 - XML 源 - Python - 如何?
【发布时间】:2020-02-05 03:49:35
【问题描述】:

我正在尝试将 XML 文件导入我的数据流代码。我看到 java 有一个内置的 XMLIo 但 Python 没有?我自己也很难理解 ParDo/DoFn 的初始步骤是什么。这是以下 XML 文件的示例。我在解析 .csv 时理解下面的管道,但我不明白如何从 XML 源开始。我是否需要手动创建 PCollection 并从那里开始?

我的目标是将每个元素作为一个元组返回。键是国家名称,后面的每个元素(在嵌套数组中)都是值。

<?xml version="1.0"?>
<data>
    <country name="Liechtenstein">
        <rank>1</rank>
        <year>2008</year>
        <gdppc>141100</gdppc>
        <neighbor name="Austria" direction="E"/>
        <neighbor name="Switzerland" direction="W"/>
    </country>
    <country name="Singapore">
        <rank>4</rank>
        <year>2011</year>
        <gdppc>59900</gdppc>
        <neighbor name="Malaysia" direction="N"/>
    </country>
    <country name="Panama">
        <rank>68</rank>
        <year>2011</year>
        <gdppc>13600</gdppc>
        <neighbor name="Costa Rica" direction="W"/>
        <neighbor name="Colombia" direction="E"/>
    </country>
</data>

def run():
   argv = [
      '--project={0}'.format(PROJECT),
      '--staging_location=gs://{0}/'.format(BUCKET),
      '--temp_location=gs://{0}/'.format(BUCKET),
      '--runner=DataflowRunner'
      #'--runner=DirectRunner'
   ]

   p = beam.Pipeline(argv=argv)

   (p
      | 'ReadFromGCS' >> beam.io.textio.ReadFromText('gs://{0}/example.csv'.format(BUCKET))
-[SNIP]-

【问题讨论】:

    标签: python xml apache-beam dataflow


    【解决方案1】:

    下面的代码将收集每个国家的信息。

    输出是一个元组列表。

    元组中的第一个元素是国家名称,第二个元素是其他国家属性的列表。

    import xml.etree.ElementTree as ET
    
    
    xml = '''<?xml version="1.0"?>
    <data>
        <country name="Liechtenstein">
            <rank>1</rank>
            <year>2008</year>
            <gdppc>141100</gdppc>
            <neighbor name="Austria" direction="E"/>
            <neighbor name="Switzerland" direction="W"/>
        </country>
        <country name="Singapore">
            <rank>4</rank>
            <year>2011</year>
            <gdppc>59900</gdppc>
            <neighbor name="Malaysia" direction="N"/>
        </country>
        <country name="Panama">
            <rank>68</rank>
            <year>2011</year>
            <gdppc>13600</gdppc>
            <neighbor name="Costa Rica" direction="W"/>
            <neighbor name="Colombia" direction="E"/>
        </country>
    </data>'''
    
    result = []
    root = ET.fromstring(xml)
    for country in root.findall('.//country'):
        result.append((country.attrib['name'],[x.text if x.text else x.attrib for x in list(country)]))
    print(result)
    

    输出

    [('Liechtenstein', ['1', '2008', '141100', {'name': 'Austria', 'direction': 'E'}, {'name': 'Switzerland','direction': 'W'}]), ('Singapore', ['4', '2011', '59900', {'name': 'Malaysia', 'direction': 'N'}]), ('Panama', ['68', '2011', '13600', {'name': 'Costa Rica', 'direction': 'W'}, {'name': 'Colombia', 'direction': 'E'}])]
    

    【讨论】:

    • 感谢 balderman!你会碰巧知道这将如何在梁管道中工作吗?如果我从包含所有国家/地区条目的元组中形成一个 PCollection(在我正在使用的生产数据集中可能是数十而不是数千),我的下一次转换是否基本上需要将此单行/元组传输到所有工作节点进行处理?
    • 我对梁不熟悉。如果您发现我的 xml 解析有用 - 请随时投票。
    • 会的!让我测试一下,并在他们应得的地方提供道具:)
    • 这确实帮助我编写了一个 DoFn 来解析 XML。这是一个很大的难题,非常感谢!!!
    猜你喜欢
    • 1970-01-01
    • 2018-03-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-05-16
    • 1970-01-01
    • 1970-01-01
    • 2010-11-01
    相关资源
    最近更新 更多