From 32adca324218e407cef8eba47a40f6e267076bfb Mon Sep 17 00:00:00 2001 From: bertber <75388552+bertber@users.noreply.github.com> Date: Sat, 5 Dec 2020 02:51:27 +0800 Subject: [PATCH 1/2] MINOR: add a description for calling resetToDatetime's function to resetByDuration Add a request that it will call resetToDatetime of description to resetByDuration of function. --- .../src/main/scala/kafka/tools/StreamsResetter.java | 13 ++----------- 1 file changed, 2 insertions(+), 11 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index d356681e11242..adaee6694510d 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -491,17 +491,8 @@ private void resetByDuration(final Consumer client, 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, timestamp); } private void resetToDatetime(final Consumer client, From 59a880349e4bd5b93ceaac713572979f2fb3e2a1 Mon Sep 17 00:00:00 2001 From: bertber <75388552+bertber@users.noreply.github.com> Date: Wed, 9 Dec 2020 22:25:00 +0800 Subject: [PATCH 2/2] remove duplicate code from resetByDuration remove local variables from resetByDuration --- core/src/main/scala/kafka/tools/StreamsResetter.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index adaee6694510d..f72a3b6905400 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -489,10 +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(); - - resetToDatetime(client, inputTopicPartitions, timestamp); + resetToDatetime(client, inputTopicPartitions, Instant.now().minus(duration).toEpochMilli()); } private void resetToDatetime(final Consumer client,