diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index d356681e11242..f72a3b6905400 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -489,19 +489,7 @@ private Map getTopicPartitionOffsetFromResetPlan(final Str private void resetByDuration(final Consumer client, final Set inputTopicPartitions, final Duration duration) { - final Instant now = Instant.now(); - final long timestamp = now.minus(duration).toEpochMilli(); - - final Map topicPartitionsAndTimes = new HashMap<>(inputTopicPartitions.size()); - for (final TopicPartition topicPartition : inputTopicPartitions) { - topicPartitionsAndTimes.put(topicPartition, timestamp); - } - - final Map topicPartitionsAndOffset = client.offsetsForTimes(topicPartitionsAndTimes); - - for (final TopicPartition topicPartition : inputTopicPartitions) { - client.seek(topicPartition, topicPartitionsAndOffset.get(topicPartition).offset()); - } + resetToDatetime(client, inputTopicPartitions, Instant.now().minus(duration).toEpochMilli()); } private void resetToDatetime(final Consumer client,