Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/layouts/shortcodes/generated/pipeline_configuration.html
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,12 @@
<td>Boolean</td>
<td>Forces Flink to register avro classes in kryo serializer.<br /><br />Important: Make sure to include the flink-avro module. Otherwise, nothing will be registered. For backward compatibility, the default value is empty to conform to the behavior of the older version. That is, always register avro with kryo, and if flink-avro is not in the class path, register a dummy serializer. In Flink-2.0, we will set the default value to true.</td>
</tr>
<tr>
<td><h5>pipeline.forward-edge.parallelism-mismatch-mode</h5></td>
<td style="word-wrap: break-word;">REBALANCE</td>
<td><p>Enum</p></td>
<td>Determines how the runtime handles a FORWARD (pointwise) edge whose producer and consumer parallelism no longer match, which can happen when the parallelism of connected operators is changed independently (for example via the AdaptiveScheduler or <code class="highlighter-rouge">pipeline.jobvertex-parallelism-overrides</code>). A FORWARD edge is only valid at equal parallelism; on a mismatch one of the following strategies is applied:<ul><li><code class="highlighter-rouge">REBALANCE</code>: replace the forward partitioner with a rebalance (round-robin) partitioner. This preserves throughput but reorders records, which is safe for append-only streams but corrupts order-sensitive (changelog) streams.</li><li><code class="highlighter-rouge">FAIL</code>: fail the job with an exception instead of silently changing the data distribution.</li><li><code class="highlighter-rouge">KEEP_FORWARD</code>: keep the forward partitioner. Record order is preserved but the records are funneled to a single consumer subtask, so the remaining consumer subtasks stay idle.</li></ul><br /><br />Possible values:<ul><li>"REBALANCE"</li><li>"FAIL"</li><li>"KEEP_FORWARD"</li></ul></td>
</tr>
<tr>
<td><h5>pipeline.generic-types</h5></td>
<td style="word-wrap: break-word;">true</td>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -305,6 +305,68 @@ public enum VertexDescriptionMode {
CASCADING
}

public static final ConfigOption<ForwardEdgeParallelismMismatchMode>
FORWARD_EDGE_PARALLELISM_MISMATCH_MODE =
key("pipeline.forward-edge.parallelism-mismatch-mode")
.enumType(ForwardEdgeParallelismMismatchMode.class)
.defaultValue(ForwardEdgeParallelismMismatchMode.REBALANCE)
.withDescription(
Description.builder()
.text(
"Determines how the runtime handles a FORWARD (pointwise) edge whose"
+ " producer and consumer parallelism no longer match, which can happen"
+ " when the parallelism of connected operators is changed"
+ " independently (for example via the AdaptiveScheduler or %s). A"
+ " FORWARD edge is only valid at equal parallelism; on a mismatch one"
+ " of the following strategies is applied:",
code(PARALLELISM_OVERRIDES.key()))
.list(
text(
"%s: replace the forward partitioner with a rebalance (round-robin)"
+ " partitioner. This preserves throughput but reorders records,"
+ " which is safe for append-only streams but corrupts"
+ " order-sensitive (changelog) streams.",
code(
ForwardEdgeParallelismMismatchMode
.REBALANCE
.name())),
text(
"%s: fail the job with an exception instead of silently changing"
+ " the data distribution.",
code(
ForwardEdgeParallelismMismatchMode
.FAIL
.name())),
text(
"%s: keep the forward partitioner. Record order is preserved but"
+ " the records are funneled to a single consumer subtask, so the"
+ " remaining consumer subtasks stay idle.",
code(
ForwardEdgeParallelismMismatchMode
.KEEP_FORWARD
.name())))
.build());

/**
* The strategy applied when a FORWARD (pointwise) edge connects a producer and consumer with
* mismatched parallelism.
*/
@PublicEvolving
public enum ForwardEdgeParallelismMismatchMode {
/**
* Replace the forward partitioner with a rebalance partitioner. Reorders records; safe only
* for append-only streams.
*/
REBALANCE,
/** Fail the job with an exception. */
FAIL,
/**
* Keep the forward partitioner. Preserves record order but funnels records to a single
* consumer subtask.
*/
KEEP_FORWARD
}

public static final ConfigOption<Boolean> VERTEX_NAME_INCLUDE_INDEX_PREFIX =
key("pipeline.vertex-name-include-index-prefix")
.booleanType()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -322,6 +322,11 @@ public ResultPartitionWriter getWriter(int index) {
throw new UnsupportedOperationException(ERROR_MSG);
}

@Override
public int getWriterConsumerParallelism(int index) {
throw new UnsupportedOperationException(ERROR_MSG);
}

@Override
public ResultPartitionWriter[] getAllWriters() {
return new ResultPartitionWriter[0];
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@

import java.io.Serializable;

import static org.apache.flink.util.Preconditions.checkArgument;
import static org.apache.flink.util.Preconditions.checkNotNull;

/**
Expand All @@ -45,14 +46,31 @@ public class ResultPartitionDeploymentDescriptor implements Serializable {

private final int maxParallelism;

/** Parallelism of the consuming job vertices, or {@link #UNKNOWN_CONSUMER_PARALLELISM}. */
private final int consumerParallelism;

/**
* The consumer parallelism is not decided yet when the producer is deployed. This only happens
* in dynamic graphs (AdaptiveBatchScheduler), where the parallelism of a consumer vertex may be
* decided after its producers are deployed. Every other partition carries the actual consumer
* parallelism.
*/
public static final int UNKNOWN_CONSUMER_PARALLELISM = -1;

public ResultPartitionDeploymentDescriptor(
PartitionDescriptor partitionDescriptor,
ShuffleDescriptor shuffleDescriptor,
int maxParallelism) {
int maxParallelism,
int consumerParallelism) {
this.partitionDescriptor = checkNotNull(partitionDescriptor);
this.shuffleDescriptor = checkNotNull(shuffleDescriptor);
KeyGroupRangeAssignment.checkParallelismPreconditions(maxParallelism);
this.maxParallelism = maxParallelism;
checkArgument(
consumerParallelism > 0 || consumerParallelism == UNKNOWN_CONSUMER_PARALLELISM,
"Invalid consumer parallelism %s.",
consumerParallelism);
this.consumerParallelism = consumerParallelism;
}

public IntermediateDataSetID getResultId() {
Expand Down Expand Up @@ -88,6 +106,16 @@ public int getMaxParallelism() {
return maxParallelism;
}

/**
* Returns the parallelism of the job vertices consuming this partition, or {@link
* #UNKNOWN_CONSUMER_PARALLELISM} if it is not decided yet. Unlike {@link
* #getNumberOfSubpartitions()}, this is the actual consumer parallelism for every distribution
* pattern.
*/
public int getConsumerParallelism() {
return consumerParallelism;
}

public ShuffleDescriptor getShuffleDescriptor() {
return shuffleDescriptor;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.apache.flink.runtime.checkpoint.TaskStateSnapshot;
import org.apache.flink.runtime.checkpoint.channel.ChannelStateWriteRequestExecutorFactory;
import org.apache.flink.runtime.checkpoint.channel.ChannelStateWriter;
import org.apache.flink.runtime.deployment.ResultPartitionDeploymentDescriptor;
import org.apache.flink.runtime.executiongraph.ExecutionAttemptID;
import org.apache.flink.runtime.externalresource.ExternalResourceInfoProvider;
import org.apache.flink.runtime.io.disk.iomanager.IOManager;
Expand Down Expand Up @@ -248,6 +249,13 @@ void acknowledgeCheckpoint(

ResultPartitionWriter getWriter(int index);

/**
* Returns the parallelism of the job vertices consuming the partition written by {@link
* #getWriter(int)}, or {@link ResultPartitionDeploymentDescriptor#UNKNOWN_CONSUMER_PARALLELISM}
* if it is not decided yet.
*/
int getWriterConsumerParallelism(int index);

ResultPartitionWriter[] getAllWriters();

IndexedInputGate getInputGate(int index);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -561,7 +561,10 @@ private static ResultPartitionDeploymentDescriptor createResultPartitionDeployme
IntermediateResultPartition partition,
ShuffleDescriptor shuffleDescriptor) {
return new ResultPartitionDeploymentDescriptor(
partitionDescriptor, shuffleDescriptor, getPartitionMaxParallelism(partition));
partitionDescriptor,
shuffleDescriptor,
getPartitionMaxParallelism(partition),
partition.getIntermediateResult().getConsumersParallelism());
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;

import static org.apache.flink.util.Preconditions.checkArgument;
import static org.apache.flink.util.Preconditions.checkNotNull;
import static org.apache.flink.util.Preconditions.checkState;

Expand Down Expand Up @@ -93,6 +94,7 @@ public class RuntimeEnvironment implements Environment {
private final Map<String, Future<Path>> distCacheEntries;

private final ResultPartitionWriter[] writers;
private final int[] writerConsumerParallelisms;
private final IndexedInputGate[] inputGates;

private final TaskEventDispatcher taskEventDispatcher;
Expand Down Expand Up @@ -145,6 +147,7 @@ public RuntimeEnvironment(
InputSplitProvider splitProvider,
Map<String, Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
int[] writerConsumerParallelisms,
IndexedInputGate[] inputGates,
TaskEventDispatcher taskEventDispatcher,
CheckpointResponder checkpointResponder,
Expand Down Expand Up @@ -177,6 +180,8 @@ public RuntimeEnvironment(
this.splitProvider = checkNotNull(splitProvider);
this.distCacheEntries = checkNotNull(distCacheEntries);
this.writers = checkNotNull(writers);
this.writerConsumerParallelisms = checkNotNull(writerConsumerParallelisms);
checkArgument(writers.length == writerConsumerParallelisms.length);
this.inputGates = checkNotNull(inputGates);
this.taskEventDispatcher = checkNotNull(taskEventDispatcher);
this.checkpointResponder = checkNotNull(checkpointResponder);
Expand Down Expand Up @@ -306,6 +311,11 @@ public ResultPartitionWriter getWriter(int index) {
return writers[index];
}

@Override
public int getWriterConsumerParallelism(int index) {
return writerConsumerParallelisms[index];
}

@Override
public ResultPartitionWriter[] getAllWriters() {
return writers;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,9 @@ public class Task

private final ResultPartitionWriter[] partitionWriters;

/** Consumer parallelism of each partition in {@link #partitionWriters}, in the same order. */
private final int[] partitionConsumerParallelisms;

private final IndexedInputGate[] inputGates;

/** Connection to the task manager. */
Expand Down Expand Up @@ -423,6 +426,10 @@ public Task(
.toArray(new ResultPartitionWriter[] {});

this.partitionWriters = resultPartitionWriters;
this.partitionConsumerParallelisms =
resultPartitionDeploymentDescriptors.stream()
.mapToInt(ResultPartitionDeploymentDescriptor::getConsumerParallelism)
.toArray();

// consumed intermediate result partitions
final IndexedInputGate[] gates =
Expand Down Expand Up @@ -731,6 +738,7 @@ private void doRun() {
inputSplitProvider,
distributedCacheEntries,
partitionWriters,
partitionConsumerParallelisms,
inputGates,
taskEventDispatcher,
checkpointResponder,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@
import org.apache.flink.configuration.CheckpointingOptions;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.NettyShuffleEnvironmentOptions;
import org.apache.flink.configuration.PipelineOptions;
import org.apache.flink.configuration.PipelineOptions.ForwardEdgeParallelismMismatchMode;
import org.apache.flink.configuration.TaskManagerOptions;
import org.apache.flink.core.execution.RecoveryClaimMode;
import org.apache.flink.core.fs.AutoCloseableRegistry;
Expand All @@ -49,6 +51,7 @@
import org.apache.flink.runtime.checkpoint.channel.RecoveryCheckpointTrigger;
import org.apache.flink.runtime.checkpoint.channel.SequentialChannelStateReader;
import org.apache.flink.runtime.checkpoint.filemerging.FileMergingSnapshotManager;
import org.apache.flink.runtime.deployment.ResultPartitionDeploymentDescriptor;
import org.apache.flink.runtime.execution.CancelTaskException;
import org.apache.flink.runtime.execution.Environment;
import org.apache.flink.runtime.io.AvailabilityProvider;
Expand Down Expand Up @@ -2046,13 +2049,59 @@ List<RecordWriter<SerializationDelegate<StreamRecord<OUT>>>> createRecordWriters

private static void replaceForwardPartitionerIfConsumerParallelismDoesNotMatch(
Environment environment, NonChainedOutput streamOutput, int outputIndex) {
if (streamOutput.getPartitioner() instanceof ForwardPartitioner
&& environment.getWriter(outputIndex).getNumberOfSubpartitions()
!= environment.getTaskInfo().getNumberOfParallelSubtasks()) {
LOG.debug(
"Replacing forward partitioner with rebalance for {}",
environment.getTaskInfo().getTaskNameWithSubtasks());
streamOutput.setPartitioner(new RebalancePartitioner<>());
final int producerParallelism = environment.getTaskInfo().getNumberOfParallelSubtasks();
// The number of subpartitions isn't the consumer parallelism on a POINTWISE edge, so the
// actual consumer parallelism is passed through the deployment descriptor.
final int consumerParallelism = environment.getWriterConsumerParallelism(outputIndex);
if (!(streamOutput.getPartitioner() instanceof ForwardPartitioner)
// An undecided consumer parallelism is not a mismatch.
|| consumerParallelism
== ResultPartitionDeploymentDescriptor.UNKNOWN_CONSUMER_PARALLELISM
|| consumerParallelism == producerParallelism) {
return;
}

final ForwardEdgeParallelismMismatchMode mode =
environment
.getJobConfiguration()
.get(PipelineOptions.FORWARD_EDGE_PARALLELISM_MISMATCH_MODE);
final String taskNameWithSubtasks = environment.getTaskInfo().getTaskNameWithSubtasks();
switch (mode) {
case REBALANCE:
LOG.debug(
"Replacing forward partitioner with rebalance for {} "
+ "(producer parallelism {} != consumer parallelism {}).",
taskNameWithSubtasks,
producerParallelism,
consumerParallelism);
streamOutput.setPartitioner(new RebalancePartitioner<>());
break;
case KEEP_FORWARD:
LOG.warn(
"Keeping forward partitioner for {} despite a parallelism mismatch "
+ "(producer parallelism {} != consumer parallelism {}). Record order "
+ "is preserved but records are funneled to a single consumer subtask, "
+ "leaving the remaining consumer subtasks idle.",
taskNameWithSubtasks,
producerParallelism,
consumerParallelism);
break;
case FAIL:
throw new FlinkRuntimeException(
String.format(
"Forward partitioning cannot be preserved across a parallelism change "
+ "for %s (producer parallelism %d != consumer parallelism %d). "
+ "Silently downgrading a FORWARD edge to a redistributing "
+ "exchange can reorder records and corrupt order-sensitive "
+ "(changelog) results. Set %s to %s or %s to allow it.",
taskNameWithSubtasks,
producerParallelism,
consumerParallelism,
PipelineOptions.FORWARD_EDGE_PARALLELISM_MISMATCH_MODE.key(),
ForwardEdgeParallelismMismatchMode.REBALANCE,
ForwardEdgeParallelismMismatchMode.KEEP_FORWARD));
default:
throw new IllegalStateException("Unhandled mode: " + mode);
}
}

Expand Down
Loading