【问题标题】:Beam on SparkRunner overwriting its own outputSparkRunner 上的 Beam 覆盖了自己的输出
【发布时间】:2019-10-22 19:04:52
【问题描述】:

我正在 SparkRunner 上运行 Beam 管道,并带有 Parquet 文件输出(尽管如果我正在执行其他 IO 输出,问题就存在)。我遇到的问题是在输出时,文件副本覆盖了它自己的输出。这是日志输出:

19/10/22 18:26:35 INFO FileBasedSink: Will copy temporary file FileResult{tempFilename=/home/hadoop/just1hour/.temp-beam-357d6916-5d8e-4519-a7a4-3852249011b5/77100cd1-04ae-441c-848f-e0d0067feeb8, shard=0, window=org.apache.beam.sdk.transforms.windowing.GlobalWindow@5316e95f,
 paneInfo=PaneInfo.NO_FIRING} to final location /home/hadoop/just1hour/output-00000-of-00001
19/10/22 18:26:35 INFO FileBasedSink: Will copy temporary file FileResult{tempFilename=/home/hadoop/just1hour/.temp-beam-357d6916-5d8e-4519-a7a4-3852249011b5/f19819df-e006-431a-8ccd-6e67af692c3e, shard=0, window=org.apache.beam.sdk.transforms.windowing.GlobalWindow@5316e95f,
 paneInfo=PaneInfo.NO_FIRING} to final location /home/hadoop/just1hour/output-00000-of-00001
19/10/22 18:26:35 INFO FileBasedSink: Will copy temporary file FileResult{tempFilename=/home/hadoop/just1hour/.temp-beam-357d6916-5d8e-4519-a7a4-3852249011b5/cb2abe0c-8cc2-4a94-ae54-97b67c5e7d20, shard=0, window=org.apache.beam.sdk.transforms.windowing.GlobalWindow@5316e95f,
 paneInfo=PaneInfo.NO_FIRING} to final location /home/hadoop/just1hour/output-00000-of-00002
19/10/22 18:26:35 INFO FileBasedSink: Will copy temporary file FileResult{tempFilename=/home/hadoop/just1hour/.temp-beam-357d6916-5d8e-4519-a7a4-3852249011b5/611d194a-4e8f-44bc-8776-4bc2c55a8f34, shard=1, window=org.apache.beam.sdk.transforms.windowing.GlobalWindow@5316e95f,
 paneInfo=PaneInfo.NO_FIRING} to final location /home/hadoop/just1hour/output-00001-of-00002

如您所见,第一个文件正在被覆盖。

我可以通过手动指定分片数等于输入文件数来解决此问题,但我想知道是否有其他配置可以解释或避免这种行为。

编辑:

这是一个批处理作业,下面是生成输出的代码:

p.apply(TextIO.read().from(input).withDelimiter("{".getBytes()))
                .apply(Filter.by((String record) -> !record.isEmpty()))
                .apply(ParDo.of(new ParseNotificationJSON())).setCoder(AvroCoder.of(SCHEMA))
                .apply("Write Parquet files",
                        FileIO.<GenericRecord>write().via(ParquetIO.sink(SCHEMA)).to(output));

        p.run().waitUntilFinish();

【问题讨论】:

  • 它是批处理还是流处理作业?此外,如果您可以提供最小的代码块以便重新生成问题,那将是很好的。
  • 在上面添加了该信息。

标签: apache-spark apache-beam parquet


【解决方案1】:

如果目标目录中的任何文件具有相同的名称,Beam 中基于 FileIO 的接收器将覆盖这些文件。有界源的默认文件命名也在文件名中使用 shard-Index 和 shard-number,因此使用 .withNumShards(0) 将使用运行器确定的分片。如果您在接收器中使用.withNumShards(0),它应该可以正常工作。

【讨论】:

  • 不,这不起作用——同样的行为仍然存在。问题似乎是每组分片似乎都不知道其他分片。在示例日志输出中,第一个分片输出到 output-00000-of-00001,然后下一个分片覆盖第一个分片。第二组分片没问题,因为它们正在写入 output-00000-of-00002 和 output-00001-of-00002。
猜你喜欢
  • 1970-01-01
  • 2011-03-15
  • 2014-11-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-06-26
  • 2019-01-18
  • 1970-01-01
相关资源
最近更新 更多