apache / apache/beam

No parallelism when using SDFBoundedSourceReader with Flink

Open
#21,222 0 comments 0 reactions 0 assignees View on GitHub
bug flink P3 runners
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
2d 2h
Merged PRs (30d)
205

Description

Background: I am using TFX pipelines with Flink as the runner for Beam (flink session cluster using [flink-on-k8s-operator](https://github.com/GoogleCloudPlatform/flink-on-k8s-operator)). The Flink cluster has 2 taskmanagers with 16 cores each, and parallelism is set to 32. TFX components call `beam.io.ReadFromTFRecord` to load data, passing in a glob file pattern. I have a dataset of TFRecords split across 160 files. When I try to run the component, processing for all 160 files ends up in a single subtask in Flink, i.e. the parallelism is effectively 1. See below images:

!https://i.imgur.com/ppba0AL.png!

!https://i.imgur.com/rSTFATn.png!

 
I have tried all manner of Beam/Flink options and different versions of Beam/Flink but the behaviour remains the same.

Furthermore, the behaviour affects anything that uses `apache_beam.io.iobase.SDFBoundedSourceReader`, e.g. `apache_beam.io.parquetio.ReadFromParquet` also has the same issue. Either I'm missing some obscure setting in my configuration, or this is a bug with the Flink runner.
 

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

Contributor guide

Open the contributing guide

Research direction

Start by reproducing the reported behavior with Flink parallelism 32, 160 TFRecord files, and Beam's SDFBoundedSourceReader. Compare ReadFromTFRecord and ReadFromParquet, then inspect the Flink runner's handling of SDFBoundedSourceReader. Done means the source work is distributed across the configured Flink subtasks rather than processed by one subtask.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
distributed-systems
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.