Apache Flink问题:自定义InputFormat仅支持并行度为1运行
Hey there! Let's break down why your custom ElasticsearchInputFormat is stuck running with parallelism 1, and walk through how to fix it.
The Root Cause
Your format extends GenericInputFormat, which defaults to generating only one input split out of the box. Flink assigns one parallel task per input split, so with just one split, you can't scale beyond parallelism 1. To enable parallel execution, you need to implement logic to split your Elasticsearch data into multiple independent chunks that can be processed in parallel.
Step-by-Step Solution
Here's how to modify your code to support parallelism:
1. Generate Multiple Input Splits
Override the createInputSplits method to create a number of splits matching your desired parallelism. For a simple test case, we'll create splits based on numeric indices, but you can later adapt this to use Elasticsearch's own shard structure for better real-world efficiency.
@Override public GenericInputSplit[] createInputSplits(int numSplits) throws IOException { GenericInputSplit[] splits = new GenericInputSplit[numSplits]; for (int i = 0; i < numSplits; i++) { // Each split gets a unique index and the total number of splits splits[i] = new GenericInputSplit(i, numSplits); } return splits; }
2. Specify Input Split Type
Tell Flink how your splits are structured by overriding getInputSplitType. For evenly distributed splits, use EQUALITY:
@Override public InputSplitSource.InputSplitType getInputSplitType() { return InputSplitSource.InputSplitType.EQUALITY; }
3. Use Split Information to Partition Data
In the open method, grab the split's index and total split count, then use this to determine which portion of the data this parallel task should process. For your sample 3-row dataset, we'll use this to partition rows across splits:
private int splitIndex; private int totalSplits; private int currentRow = 0; // Sample test data private final Row[] sampleRows = { Row.of("row1", 1), Row.of("row2", 2), Row.of("row3", 3) }; @Override public void open(GenericInputSplit split) throws IOException { super.open(split); this.splitIndex = split.getSplitNumber(); this.totalSplits = split.getTotalNumberOfSplits(); }
4. Adjust Record Fetching Logic
Update nextRecord and reachedEnd to only return rows assigned to the current split. We'll use modulo to distribute rows evenly across splits:
@Override public boolean reachedEnd() throws IOException { return currentRow >= sampleRows.length; } @Override public Row nextRecord(Row reuse) throws IOException { // Skip rows that don't belong to this split while (currentRow < sampleRows.length) { Row row = sampleRows[currentRow]; currentRow++; // Assign row to split using row index modulo total splits if ((currentRow - 1) % totalSplits == splitIndex) { reuse.copy(row); return reuse; } } return null; }
Elasticsearch-Specific Optimization
For real-world usage, instead of manual row partitioning, you can:
- Fetch Elasticsearch's cluster shard information via its API in
createInputSplits - Map each ES shard to a Flink
InputSplit - In
open, have each task query only the assigned ES shard
This aligns Flink's parallelism with Elasticsearch's native shard structure, maximizing efficiency and avoiding overloading the ES cluster.
Testing the Fix
Once you've updated the code, set your job's parallelism (e.g., env.setParallelism(3)), and you'll see each parallel task processing its own portion of the data.
内容的提问来源于stack exchange,提问作者Mustafa

