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);
   }
 }

Reply via email to