【发布时间】: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