【问题标题】:Apache-Beam + Python: Writing JSON (or dictionaries) strings to output fileApache-Beam + Python:将 JSON(或字典)字符串写入输出文件
【发布时间】:2017-07-21 16:08:49
【问题描述】:

我正在尝试使用 Beam 管道将 SequenceMatcher 函数应用于大量单词。除了 WriteToText 部分,我(希望)已经弄清楚了一切。

我已经定义了一个自定义 ParDo(这里称为 ProcessDataDoFn),它接受 main_input 和 side_input,处理它们并输出像这样的字典

{u'key': (u'string', float)}

我的管道很简单

class ProcessDataDoFn(beam.DoFn):
    def process(self, element, side_input):

    ... Series of operations ...

    return output_dictionary

with beam.Pipeline(options=options) as p:

    # Main input
    main_input = p | 'ReadMainInput' >> beam.io.Read(
        beam.io.BigQuerySource(
            query=CUSTOM_SQL,
            use_standard_sql=True
        ))

    # Side input
    side_input = p | 'ReadSideInput' >> beam.io.Read(
        beam.io.BigQuerySource(
            project=PROJECT_ID,
            dataset=DATASET,
            table=TABLE
        ))

    output = (
        main_input
        | 'ProcessData' >> beam.ParDo(
            ProcessDataDoFn(),
            side_input=beam.pvalue.AsList(side_input))
        | 'WriteOutput' >> beam.io.WriteToText(GCS_BUCKET)
    )

现在的问题是,如果我这样离开管道,它只会输出 output_dictionary 的键。如果我将 ProcessDataDoFn 的返回更改为 json.dumps(ouput_dictionary),则 Json 写入正确,但像这样

{
'
k
e
y
'

:

[
'
s
t
r
i
n
g
'

,

f
l
o
a
t
]

如何正确输出结果?

【问题讨论】:

  • 在您的代码中,该类被声明为ProcessData,那么当您在管道中使用它时,它就是ProcessDataDoFn。我确定这只是问题中的一个错字,但有助于纠正它。
  • 谢谢,现在应该修好了。

标签: python json dictionary google-cloud-dataflow apache-beam


【解决方案1】:

您的输出看起来像这样是不寻常的。 json.dumps 应该在一行中打印 json,并且应该逐行输出到文件。

也许有更简洁的代码,您可以添加一个额外的地图操作,以您需要的方式进行格式化。像这样:

output = (
  main_input
  | 'ProcessData' >> beam.ParDo(
        ProcessDataDoFn(),
        side_input=beam.pvalue.AsList(side_input))
  | 'FormatOutput' >> beam.Map(json.dumps)
  | 'WriteOutput' >> beam.io.WriteToText(GCS_BUCKET)
)

【讨论】:

  • 是否有理由推荐将json.dumps 函数映射到pcoll 中的元素,而不是使用ParDo 来应用DoFn?似乎 DoFn 具有指标和其他好处,或者您认为对于像这样的简单转换案例来说开销太大?
  • A Map 被翻译成 ParDo,所以任何适合您的用例都可以。如果你想要额外的功能,去ParDo。 :)
【解决方案2】:

我实际上部分解决了这个问题。

我编写的 ParDoFn 要么返回字典,要么返回 JSON 格式的字符串。在这两种情况下,当 Beam 尝试对所述输入执行某些操作时,就会出现问题。如果所述 PCollection 是字典,Beam 似乎会遍历给定的 PCollection,它只获取它的键,如果所述 PCollection 是字符串,它会遍历所有字符(这就是 JSON 输出如此奇怪的原因)。我发现解决方案相当简单:将字典或字符串封装在列表中。 JSON 格式化部分可以在 ParDoFn 级别完成,也可以通过像您展示的那样的转换来完成。

【讨论】:

  • 如 Python API here 中所述:“请注意,DoFn 必须为输入 PCollection 的每个元素返回一个可迭代对象。一个简单的方法是在过程中使用 yield 关键字方法。”
  • 如果可以的话,我很想见到你的ProcessDataDoFn()。您是否最终对列表中的每个字典都使用了yield
  • 不幸的是,它是专有代码,但是,yield 为我解决了它。
  • 字符串是可迭代的,但随后切片 i 和列表函数成为可能,这是 python 的一个奇怪的怪癖。我猜 PCollection 只是期待一个可迭代的,并且可能有一个用例将字符串转换为 PCollection 以逐个字符处理。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-06-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-08-26
相关资源
最近更新 更多