WencongLiu commented on code in PR #23179:
URL: https://github.com/apache/flink/pull/23179#discussion_r1294530493


##########
flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/hybrid/tiered/tier/disk/DiskIOScheduler.java:
##########
@@ -236,18 +236,18 @@ private int readBuffersFromFile() {
     }
 
     private List<ScheduledSubpartitionReader> sortScheduledReaders() {
+        List<ScheduledSubpartitionReader> scheduledReaders;
         synchronized (lock) {
             if (isReleased) {
                 return new ArrayList<>();
             }
-            List<ScheduledSubpartitionReader> scheduledReaders = new 
ArrayList<>();
-            for (ScheduledSubpartitionReader reader : 
allScheduledReaders.values()) {
-                reader.prepareForScheduling();
-                scheduledReaders.add(reader);
-            }
-            Collections.sort(scheduledReaders);
-            return scheduledReaders;
+            scheduledReaders = new ArrayList<>(allScheduledReaders.values());
+        }
+        for (ScheduledSubpartitionReader reader : scheduledReaders) {

Review Comment:
   I've add a test in `DiskIOScheduler`.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/hybrid/HsSubpartitionConsumerMemoryDataManager.java:
##########
@@ -173,7 +179,19 @@ public void releaseDataView() {
     @GuardedBy("consumerLock")
     private void trimHeadingReleasedBuffers() {
         while (!unConsumedBuffers.isEmpty() && 
unConsumedBuffers.peekFirst().isReleased()) {
-            unConsumedBuffers.removeFirst();
+            decreaseBacklog(unConsumedBuffers.removeFirst().getBuffer());
+        }
+    }
+
+    private void increaseBacklog(Buffer buffer) {
+        if (buffer.isBuffer()) {
+            backlog.getAndIncrement();
+        }
+    }
+
+    private void decreaseBacklog(Buffer buffer) {

Review Comment:
   Fixed.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/hybrid/HsSubpartitionConsumerMemoryDataManager.java:
##########
@@ -173,7 +179,19 @@ public void releaseDataView() {
     @GuardedBy("consumerLock")
     private void trimHeadingReleasedBuffers() {
         while (!unConsumedBuffers.isEmpty() && 
unConsumedBuffers.peekFirst().isReleased()) {
-            unConsumedBuffers.removeFirst();
+            decreaseBacklog(unConsumedBuffers.removeFirst().getBuffer());
+        }
+    }
+
+    private void increaseBacklog(Buffer buffer) {

Review Comment:
   Fixed.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/hybrid/HsSubpartitionFileReaderImpl.java:
##########
@@ -72,6 +73,9 @@ public class HsSubpartitionFileReaderImpl implements 
HsSubpartitionFileReader {
 
     private final Consumer<HsSubpartitionFileReader> fileReaderReleaser;
 
+    @GuardedBy("lock")

Review Comment:
   Removed.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to