apache / apache/beam

FileIO can produce duplicates in output files

Open
#21,082 0 comments 0 reactions 0 assignees View on GitHub
beam-model bug core files flink io java P3 runners spark
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

FileIO can produce duplicates in output files - depending on a runner.

Concrete example for Spark when executing as batch:

When using FileIO with specific number of shards, it will use default sharding function which is a round robin shard assignment with random seed. In multistage pipeline, data between stages are hold by shuffle service until downstream stage request it for further computations. If shuffle results computed with this seeded shard function are lost - e.g. shuffle service fails because of HW error - then Spark will attempt to recover data by computing them again from source data. As a result of a random seed sharding, this will assign different shard - and therefore key to the element.

More details are discussed in this thread:
https://lists.apache.org/thread.html/r5e91d1996479defbf5e896dca3cf237ee2d9b59396cb3c4edf619df1%40%3Cdev.beam.apache.org%3E

Imported from Jira [BEAM-12654](https://issues.apache.org/jira/browse/BEAM-12654). Original Jira may contain additional context.
Reported by: jvilcek.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.