This commit is contained in:
mcdull_zhang 2022-09-22 17:42:25 +08:00 committed by GitHub
parent 8e8645bb01
commit cba7dc6f71
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
2 changed files with 4 additions and 4 deletions

View File

@ -59,7 +59,7 @@ class RssShuffleReader[K, C](
readMetrics.incFetchWaitTime(time)
}
val recordIter = (startPartition until endPartition).map(partitionId => {
val recordIter = (startPartition until endPartition).iterator.map(partitionId => {
if (handle.numMaps > 0) {
val start = System.currentTimeMillis()
val inputStream = essShuffleClient.readPartition(
@ -77,7 +77,7 @@ class RssShuffleReader[K, C](
} else {
RssInputStream.empty()
}
}).toIterator.flatMap(
}).flatMap(
serializerInstance.deserializeStream(_).asKeyValueIterator)
val metricIter = CompletionIterator[(Any, Any), Iterator[(Any, Any)]](

View File

@ -60,7 +60,7 @@ class RssShuffleReader[K, C](
metrics.incFetchWaitTime(time)
}
val recordIter = (startPartition until endPartition).map(partitionId => {
val recordIter = (startPartition until endPartition).iterator.map(partitionId => {
if (handle.numMappers > 0) {
val start = System.currentTimeMillis()
val inputStream = rssShuffleClient.readPartition(
@ -79,7 +79,7 @@ class RssShuffleReader[K, C](
} else {
RssInputStream.empty()
}
}).toIterator.flatMap(
}).flatMap(
serializerInstance.deserializeStream(_).asKeyValueIterator)
val metricIter = CompletionIterator[(Any, Any), Iterator[(Any, Any)]](