From e459166931eff4b1f733fbf658f9d3a61a4efd06 Mon Sep 17 00:00:00 2001 From: Raghav Aggarwal Date: Wed, 5 Aug 2026 23:26:09 +0530 Subject: [PATCH 1/2] TEZ-4466: Measure exact stream read time while fetching --- .../tez/http/MeasuredDataInputStream.java | 79 +++++++++++++++++++ .../library/api/TezRuntimeConfiguration.java | 7 ++ .../library/common/shuffle/Fetcher.java | 5 ++ .../orderedgrouped/FetcherOrderedGrouped.java | 6 ++ 4 files changed, 97 insertions(+) create mode 100644 tez-runtime-library/src/main/java/org/apache/tez/http/MeasuredDataInputStream.java diff --git a/tez-runtime-library/src/main/java/org/apache/tez/http/MeasuredDataInputStream.java b/tez-runtime-library/src/main/java/org/apache/tez/http/MeasuredDataInputStream.java new file mode 100644 index 0000000000..46efb0df31 --- /dev/null +++ b/tez-runtime-library/src/main/java/org/apache/tez/http/MeasuredDataInputStream.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.tez.http; + +import java.io.DataInputStream; +import java.io.FilterInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.util.concurrent.TimeUnit; + +public class MeasuredDataInputStream extends DataInputStream { + + private final MeasuredInputStream measuredIn; + + private MeasuredDataInputStream(MeasuredInputStream measuredIn) { + super(measuredIn); + this.measuredIn = measuredIn; + } + + public MeasuredDataInputStream(InputStream in) { + this(new MeasuredInputStream(in)); + } + + public long getElapsedTimeMs() { + return measuredIn.getElapsedTimeMs(); + } + + private static class MeasuredInputStream extends FilterInputStream { + private long elapsedTimeNanos = 0; + + public MeasuredInputStream(InputStream in) { + super(in); + } + + @Override + public int read() throws IOException { + long start = System.nanoTime(); + int ret = super.read(); + elapsedTimeNanos += (System.nanoTime() - start); + return ret; + } + + @Override + public int read(byte[] b) throws IOException { + long start = System.nanoTime(); + int ret = super.read(b); + elapsedTimeNanos += (System.nanoTime() - start); + return ret; + } + + @Override + public int read(byte[] b, int off, int len) throws IOException { + long start = System.nanoTime(); + int ret = super.read(b, off, len); + elapsedTimeNanos += (System.nanoTime() - start); + return ret; + } + + public long getElapsedTimeMs() { + return TimeUnit.NANOSECONDS.toMillis(elapsedTimeNanos); + } + } +} diff --git a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/api/TezRuntimeConfiguration.java b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/api/TezRuntimeConfiguration.java index 569cde6367..1b70c55f3c 100644 --- a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/api/TezRuntimeConfiguration.java +++ b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/api/TezRuntimeConfiguration.java @@ -416,6 +416,13 @@ private TezRuntimeConfiguration() {} public static final float TEZ_RUNTIME_SHUFFLE_FETCH_BUFFER_PERCENT_DEFAULT = 0.90f; + /** + * Enables measuring network IO time in shuffle fetchers. + */ + @ConfigurationProperty(type = "boolean") + public static final String TEZ_RUNTIME_SHUFFLE_MEASURE_IO_TIME = TEZ_RUNTIME_PREFIX + "shuffle.measure.io.time"; + public static final boolean TEZ_RUNTIME_SHUFFLE_MEASURE_IO_TIME_DEFAULT = false; + /** * Enables fetch failures by a configuration. Should be used for testing only. */ diff --git a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java index f31140a316..dd7b7efd94 100644 --- a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java +++ b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java @@ -55,6 +55,7 @@ import org.apache.tez.dag.api.TezUncheckedException; import org.apache.tez.http.BaseHttpConnection; import org.apache.tez.http.HttpConnectionParams; +import org.apache.tez.http.MeasuredDataInputStream; import org.apache.tez.runtime.api.InputContext; import org.apache.tez.runtime.library.api.TezRuntimeConfiguration; import org.apache.tez.runtime.library.common.CompositeInputAttemptIdentifier; @@ -565,6 +566,10 @@ private HostFetchResult setupConnection(Collection attem protected void setupConnectionInternal(String host, Collection attempts) throws IOException, InterruptedException { input = httpConnection.getInputStream(); + if (conf.getBoolean(TezRuntimeConfiguration.TEZ_RUNTIME_SHUFFLE_MEASURE_IO_TIME, + TezRuntimeConfiguration.TEZ_RUNTIME_SHUFFLE_MEASURE_IO_TIME_DEFAULT)) { + input = new MeasuredDataInputStream(input); + } httpConnection.validate(); } diff --git a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java index 7f8f98f3ea..2ce59928f5 100644 --- a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java +++ b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java @@ -44,7 +44,9 @@ import org.apache.tez.common.security.JobTokenSecretManager; import org.apache.tez.http.BaseHttpConnection; import org.apache.tez.http.HttpConnectionParams; +import org.apache.tez.http.MeasuredDataInputStream; import org.apache.tez.runtime.api.InputContext; +import org.apache.tez.runtime.library.api.TezRuntimeConfiguration; import org.apache.tez.runtime.library.common.Constants; import org.apache.tez.runtime.library.common.InputAttemptIdentifier; import org.apache.tez.runtime.library.common.shuffle.InputAttemptFetchFailure; @@ -392,6 +394,10 @@ boolean setupConnection(MapHost host, Collection attempt protected void setupConnectionInternal(MapHost host, Collection attempts) throws IOException, InterruptedException { input = httpConnection.getInputStream(); + if (conf.getBoolean(TezRuntimeConfiguration.TEZ_RUNTIME_SHUFFLE_MEASURE_IO_TIME, + TezRuntimeConfiguration.TEZ_RUNTIME_SHUFFLE_MEASURE_IO_TIME_DEFAULT)) { + input = new MeasuredDataInputStream(input); + } httpConnection.validate(); } From ae7659baf2cfdf5ad207967854453cc09a954d7d Mon Sep 17 00:00:00 2001 From: Raghav Aggarwal Date: Wed, 5 Aug 2026 23:31:58 +0530 Subject: [PATCH 2/2] TEZ-4748: Expose SHUFFLE_IO_TIME_MILLISECONDS via TaskCounter --- .../java/org/apache/tez/common/counters/TaskCounter.java | 5 +++++ .../tez/runtime/library/common/shuffle/Fetcher.java | 8 ++++++++ .../shuffle/orderedgrouped/FetcherOrderedGrouped.java | 6 ++++++ 3 files changed, 19 insertions(+) diff --git a/tez-api/src/main/java/org/apache/tez/common/counters/TaskCounter.java b/tez-api/src/main/java/org/apache/tez/common/counters/TaskCounter.java index 56cafbc162..cabccd2ea1 100644 --- a/tez-api/src/main/java/org/apache/tez/common/counters/TaskCounter.java +++ b/tez-api/src/main/java/org/apache/tez/common/counters/TaskCounter.java @@ -189,6 +189,11 @@ public enum TaskCounter { */ SHUFFLE_BYTES_DISK_DIRECT, + /** + * Time spent waiting on network I/O during shuffle. Represented in milliseconds. + */ + SHUFFLE_IO_TIME_MILLISECONDS, + /** * Number of Memory to Disk merges performed during sort-merge. * Used by ShuffledMergedInput diff --git a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java index dd7b7efd94..d4a9b317f2 100644 --- a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java +++ b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/Fetcher.java @@ -51,6 +51,8 @@ import org.apache.tez.common.CallableWithNdc; import org.apache.tez.common.Preconditions; import org.apache.tez.common.TezUtilsInternal; +import org.apache.tez.common.counters.TaskCounter; +import org.apache.tez.common.counters.TezCounter; import org.apache.tez.common.security.JobTokenSecretManager; import org.apache.tez.dag.api.TezUncheckedException; import org.apache.tez.http.BaseHttpConnection; @@ -178,6 +180,7 @@ public String getHost() { BaseHttpConnection httpConnection; private HttpConnectionParams httpConnectionParams; + private final TezCounter ioTimeCounter; private final boolean localDiskFetchEnabled; private final boolean sharedFetchEnabled; @@ -220,6 +223,8 @@ protected Fetcher(FetcherCallback fetcherCallback, HttpConnectionParams params, this.localDiskFetchEnabled = localDiskFetchEnabled; this.sharedFetchEnabled = sharedFetchEnabled; + this.ioTimeCounter = inputContext.getCounters().findCounter(TaskCounter.SHUFFLE_IO_TIME_MILLISECONDS); + this.fetcherIdentifier = fetcherIdGen.getAndIncrement(); String sourceDestNameTrimmed = TezUtilsInternal.cleanVertexName(inputContext.getSourceVertexName()) + " -> " @@ -818,6 +823,9 @@ private void shutdownInternal(boolean disconnect) { synchronized (isShutDown) { try { if (httpConnection != null) { + if (input instanceof MeasuredDataInputStream && ioTimeCounter != null) { + ioTimeCounter.increment(((MeasuredDataInputStream) input).getElapsedTimeMs()); + } httpConnection.cleanup(disconnect); } } catch (IOException e) { diff --git a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java index 2ce59928f5..5d7cc758f8 100644 --- a/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java +++ b/tez-runtime-library/src/main/java/org/apache/tez/runtime/library/common/shuffle/orderedgrouped/FetcherOrderedGrouped.java @@ -40,6 +40,7 @@ import org.apache.tez.common.CallableWithNdc; import org.apache.tez.common.TezRuntimeFrameworkConfigs; import org.apache.tez.common.TezUtilsInternal; +import org.apache.tez.common.counters.TaskCounter; import org.apache.tez.common.counters.TezCounter; import org.apache.tez.common.security.JobTokenSecretManager; import org.apache.tez.http.BaseHttpConnection; @@ -77,6 +78,7 @@ class FetcherOrderedGrouped extends CallableWithNdc { private final TezCounter wrongLengthErrs; private final TezCounter badIdErrs; private final TezCounter wrongReduceErrs; + private final TezCounter ioTimeCounter; private final FetchedInputAllocatorOrderedGrouped allocator; private final ShuffleScheduler scheduler; private final ExceptionReporter exceptionReporter; @@ -153,6 +155,7 @@ public FetcherOrderedGrouped(HttpConnectionParams httpConnectionParams, this.badIdErrs = badIdErrsCounter; this.connectionErrs = connectionErrsCounter; this.wrongReduceErrs = wrongReduceErrsCounter; + this.ioTimeCounter = inputContext.getCounters().findCounter(TaskCounter.SHUFFLE_IO_TIME_MILLISECONDS); this.applicationId = inputContext.getApplicationId().toString(); this.dagId = inputContext.getDagIdentifier(); @@ -229,6 +232,9 @@ private void cleanupCurrentConnection(boolean disconnect) { synchronized (cleanupLock) { try { if (httpConnection != null) { + if (input instanceof MeasuredDataInputStream && ioTimeCounter != null) { + ioTimeCounter.increment(((MeasuredDataInputStream) input).getElapsedTimeMs()); + } httpConnection.cleanup(disconnect); httpConnection = null; }