This is an automated email from the ASF dual-hosted git repository.
zhongqiangchen pushed a commit to branch branch-0.3
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/branch-0.3 by this push:
new 7e45acb34 [CELEBORN-717][FLINK][FOLLOWUP] Fix ResultPartition lost
numBytesOut/numBuffersOut metrics
7e45acb34 is described below
commit 7e45acb34d8bf89ea853594085463443ed6f7d6a
Author: Shuang <[email protected]>
AuthorDate: Tue Jun 27 21:47:41 2023 +0800
[CELEBORN-717][FLINK][FOLLOWUP] Fix ResultPartition lost
numBytesOut/numBuffersOut metrics
### What changes were proposed in this pull request?
Metics update logic need align with Flink 1.17/1.15
### Why are the changes needed?
See [1626](https://github.com/apache/incubator-celeborn/pull/1626) And
metics update logic need align with Flink 1.17/1.15
### Does this PR introduce _any_ user-facing change?
No
### How was this patch tested?
Tpcds Manual
Closes #1631 from RexXiong/CELEBORN-717-FOLLOWUP.
Authored-by: Shuang <[email protected]>
Signed-off-by: zhongqiang.czq <[email protected]>
(cherry picked from commit fe2f76dba6ad8c6c130a23527ba9c94c225909e5)
Signed-off-by: zhongqiang.czq <[email protected]>
---
.../flink/RemoteShuffleInputGateDelegation.java | 4 +++-
.../RemoteShuffleResultPartitionDelegation.java | 27 +++++++---------------
.../celeborn/plugin/flink/buffer/SortBuffer.java | 2 +-
.../plugin/flink/RemoteShuffleResultPartition.java | 12 +++++++++-
.../plugin/flink/RemoteShuffleResultPartition.java | 16 ++++++++-----
.../plugin/flink/RemoteShuffleResultPartition.java | 20 +++++++++++-----
6 files changed, 47 insertions(+), 34 deletions(-)
diff --git
a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java
b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java
index 823cad622..d4956f692 100644
---
a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java
+++
b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java
@@ -442,7 +442,9 @@ public class RemoteShuffleInputGateDelegation {
numUnconsumedSubpartitions--;
// not the real end.
if (numSubPartitionsNotConsumed[channelInfo.getInputChannelIdx()] != 0) {
- return Optional.empty();
+ LOG.debug(
+ "numSubPartitionsNotConsumed: {}",
+ numSubPartitionsNotConsumed[channelInfo.getInputChannelIdx()]);
} else {
// the real end.
bufferReaders.get(channelIndexToReaderIndex[channelInfo.getInputChannelIdx()]).close();
diff --git
a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionDelegation.java
b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionDelegation.java
index 6298cd2ae..703fcf10f 100644
---
a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionDelegation.java
+++
b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionDelegation.java
@@ -20,11 +20,11 @@ package org.apache.celeborn.plugin.flink;
import java.io.IOException;
import java.nio.ByteBuffer;
+import java.util.function.BiConsumer;
import java.util.function.Function;
import org.apache.flink.annotation.VisibleForTesting;
import org.apache.flink.core.memory.MemorySegment;
-import org.apache.flink.metrics.Counter;
import org.apache.flink.runtime.io.network.buffer.Buffer;
import org.apache.flink.runtime.io.network.buffer.BufferCompressor;
import org.apache.flink.runtime.io.network.buffer.BufferPool;
@@ -58,25 +58,22 @@ public class RemoteShuffleResultPartitionDelegation {
/** Whether notifyEndOfData has been called or not. */
private boolean endOfDataNotified;
- private Counter numBuffersOut;
- private Counter numBytesOut;
private int numSubpartitions;
private BufferPool bufferPool;
private BufferCompressor bufferCompressor;
private Function<Buffer, Boolean> canBeCompressed;
private Runnable checkProducerState;
+ private BiConsumer<SortBuffer.BufferWithChannel, Boolean> statisticsConsumer;
public RemoteShuffleResultPartitionDelegation(
int networkBufferSize,
RemoteShuffleOutputGate outputGate,
- Counter numBuffersOut,
- Counter numBytesOut,
+ BiConsumer<SortBuffer.BufferWithChannel, Boolean> statisticsConsumer,
int numSubpartitions) {
this.networkBufferSize = networkBufferSize;
this.outputGate = outputGate;
- this.numBuffersOut = numBuffersOut;
- this.numBytesOut = numBytesOut;
this.numSubpartitions = numSubpartitions;
+ this.statisticsConsumer = statisticsConsumer;
}
public void setup(
@@ -185,7 +182,7 @@ public class RemoteShuffleResultPartitionDelegation {
Buffer buffer = bufferWithChannel.getBuffer();
int subpartitionIndex = bufferWithChannel.getChannelIndex();
- updateStatistics(bufferWithChannel.getBuffer());
+ statisticsConsumer.accept(bufferWithChannel, isBroadcast);
writeCompressedBufferIfPossible(buffer, subpartitionIndex);
}
outputGate.regionFinish();
@@ -220,11 +217,6 @@ public class RemoteShuffleResultPartitionDelegation {
outputGate.write(buffer, targetSubpartition);
}
- public void updateStatistics(Buffer buffer) {
- numBuffersOut.inc();
- numBytesOut.inc(buffer.readableBytes() - BufferUtils.HEADER_LENGTH);
- }
-
/** Spills the large record into {@link RemoteShuffleOutputGate}. */
public void writeLargeRecord(
ByteBuffer record, int targetSubpartition, Buffer.DataType dataType,
boolean isBroadcast)
@@ -242,7 +234,9 @@ public class RemoteShuffleResultPartitionDelegation {
dataType,
toCopy + BufferUtils.HEADER_LENGTH);
- updateStatistics(buffer);
+ SortBuffer.BufferWithChannel bufferWithChannel =
+ new SortBuffer.BufferWithChannel(buffer, targetSubpartition);
+ statisticsConsumer.accept(bufferWithChannel, isBroadcast);
writeCompressedBufferIfPossible(buffer, targetSubpartition);
}
outputGate.regionFinish();
@@ -332,9 +326,4 @@ public class RemoteShuffleResultPartitionDelegation {
public void setEndOfDataNotified(boolean endOfDataNotified) {
this.endOfDataNotified = endOfDataNotified;
}
-
- public void setMetricCounters(Counter numBytesOut, Counter numBuffersOut) {
- this.numBytesOut = numBytesOut;
- this.numBuffersOut = numBuffersOut;
- }
}
diff --git
a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/buffer/SortBuffer.java
b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/buffer/SortBuffer.java
index d92314abb..19f979992 100644
---
a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/buffer/SortBuffer.java
+++
b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/buffer/SortBuffer.java
@@ -74,7 +74,7 @@ public interface SortBuffer {
private final int channelIndex;
- BufferWithChannel(Buffer buffer, int channelIndex) {
+ public BufferWithChannel(Buffer buffer, int channelIndex) {
this.buffer = checkNotNull(buffer);
this.channelIndex = channelIndex;
}
diff --git
a/client-flink/flink-1.14/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
b/client-flink/flink-1.14/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
index 127b1051d..bae30ca4f 100644
---
a/client-flink/flink-1.14/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
+++
b/client-flink/flink-1.14/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
@@ -81,7 +81,10 @@ public class RemoteShuffleResultPartition extends
ResultPartition {
delegation =
new RemoteShuffleResultPartitionDelegation(
- networkBufferSize, outputGate, numBuffersOut, numBytesOut,
numSubpartitions);
+ networkBufferSize,
+ outputGate,
+ (bufferWithChannel, isBroadcast) ->
updateStatistics(bufferWithChannel, isBroadcast),
+ numSubpartitions);
}
@Override
@@ -195,4 +198,11 @@ public class RemoteShuffleResultPartition extends
ResultPartition {
public RemoteShuffleResultPartitionDelegation getDelegation() {
return delegation;
}
+
+ public void updateStatistics(
+ SortBuffer.BufferWithChannel bufferWithChannel, boolean isBroadcast) {
+ numBuffersOut.inc(isBroadcast ? numSubpartitions : 1);
+ long readableBytes = bufferWithChannel.getBuffer().readableBytes() -
BufferUtils.HEADER_LENGTH;
+ numBytesOut.inc(isBroadcast ? readableBytes * numSubpartitions :
readableBytes);
+ }
}
diff --git
a/client-flink/flink-1.15/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
b/client-flink/flink-1.15/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
index 6aa87c0f1..63106a229 100644
---
a/client-flink/flink-1.15/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
+++
b/client-flink/flink-1.15/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
@@ -34,7 +34,6 @@ import
org.apache.flink.runtime.io.network.buffer.Buffer.DataType;
import org.apache.flink.runtime.io.network.buffer.BufferCompressor;
import org.apache.flink.runtime.io.network.buffer.BufferPool;
import org.apache.flink.runtime.io.network.partition.*;
-import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
import org.apache.flink.util.function.SupplierWithException;
import org.apache.celeborn.plugin.flink.buffer.SortBuffer;
@@ -77,7 +76,10 @@ public class RemoteShuffleResultPartition extends
ResultPartition {
delegation =
new RemoteShuffleResultPartitionDelegation(
- networkBufferSize, outputGate, numBuffersOut, numBytesOut,
numSubpartitions);
+ networkBufferSize,
+ outputGate,
+ (bufferWithChannel, isBroadcast) ->
updateStatistics(bufferWithChannel, isBroadcast),
+ numSubpartitions);
}
@Override
@@ -197,9 +199,11 @@ public class RemoteShuffleResultPartition extends
ResultPartition {
return delegation;
}
- @Override
- public void setMetricGroup(TaskIOMetricGroup metrics) {
- super.setMetricGroup(metrics);
- this.delegation.setMetricCounters(numBytesOut, numBuffersOut);
+ public void updateStatistics(
+ SortBuffer.BufferWithChannel bufferWithChannel, boolean isBroadcast) {
+ numBuffersOut.inc(isBroadcast ? numSubpartitions : 1);
+ long readableBytes = bufferWithChannel.getBuffer().readableBytes() -
BufferUtils.HEADER_LENGTH;
+ numBytesProduced.inc(readableBytes);
+ numBytesOut.inc(isBroadcast ? readableBytes * numSubpartitions :
readableBytes);
}
}
diff --git
a/client-flink/flink-1.17/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
b/client-flink/flink-1.17/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
index 10b8bbe23..820e456db 100644
---
a/client-flink/flink-1.17/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
+++
b/client-flink/flink-1.17/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartition.java
@@ -35,7 +35,6 @@ import
org.apache.flink.runtime.io.network.buffer.Buffer.DataType;
import org.apache.flink.runtime.io.network.buffer.BufferCompressor;
import org.apache.flink.runtime.io.network.buffer.BufferPool;
import org.apache.flink.runtime.io.network.partition.*;
-import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
import org.apache.flink.util.function.SupplierWithException;
import org.apache.celeborn.plugin.flink.buffer.SortBuffer;
@@ -78,7 +77,10 @@ public class RemoteShuffleResultPartition extends
ResultPartition {
delegation =
new RemoteShuffleResultPartitionDelegation(
- networkBufferSize, outputGate, numBuffersOut, numBytesOut,
numSubpartitions);
+ networkBufferSize,
+ outputGate,
+ (bufferWithChannel, isBroadcast) ->
updateStatistics(bufferWithChannel, isBroadcast),
+ numSubpartitions);
}
@Override
@@ -209,9 +211,15 @@ public class RemoteShuffleResultPartition extends
ResultPartition {
return delegation;
}
- @Override
- public void setMetricGroup(TaskIOMetricGroup metrics) {
- super.setMetricGroup(metrics);
- this.delegation.setMetricCounters(numBytesOut, numBuffersOut);
+ public void updateStatistics(
+ SortBuffer.BufferWithChannel bufferWithChannel, boolean isBroadcast) {
+ numBuffersOut.inc(isBroadcast ? numSubpartitions : 1);
+ long readableBytes = bufferWithChannel.getBuffer().readableBytes() -
BufferUtils.HEADER_LENGTH;
+ if (isBroadcast) {
+ resultPartitionBytes.incAll(readableBytes);
+ } else {
+ resultPartitionBytes.inc(bufferWithChannel.getChannelIndex(),
readableBytes);
+ }
+ numBytesOut.inc(isBroadcast ? readableBytes * numSubpartitions :
readableBytes);
}
}