From 7ae63591c428e09d7c94d6c53c86b1f97926a34f Mon Sep 17 00:00:00 2001 From: Srinivas Reddy Date: Thu, 13 Dec 2018 09:23:22 +0800 Subject: [PATCH 1/4] Switch to lambda when ever possible instead of old anonymous way --- .../kafka/tools/ClientCompatibilityTest.java | 73 +++++------- .../org/apache/kafka/tools/ToolsUtils.java | 7 +- .../tools/TransactionalMessageCopier.java | 19 ++-- .../kafka/tools/VerifiableConsumer.java | 7 +- .../kafka/tools/VerifiableLog4jAppender.java | 11 +- .../kafka/tools/VerifiableProducer.java | 23 ++-- .../org/apache/kafka/trogdor/agent/Agent.java | 19 ++-- .../kafka/trogdor/agent/WorkerManager.java | 25 ++--- .../trogdor/coordinator/Coordinator.java | 19 ++-- .../kafka/trogdor/rest/JsonRestServer.java | 27 ++--- .../workload/ConnectionStressSpec.java | 7 +- .../trogdor/workload/ProduceBenchSpec.java | 7 +- .../trogdor/workload/RoundTripWorker.java | 17 ++- .../workload/RoundTripWorkloadSpec.java | 7 +- .../apache/kafka/trogdor/agent/AgentTest.java | 4 +- .../kafka/trogdor/common/ExpectedTasks.java | 106 +++++++++--------- .../trogdor/common/MiniTrogdorCluster.java | 35 +++--- .../trogdor/coordinator/CoordinatorTest.java | 10 +- .../kafka/trogdor/task/SampleTaskWorker.java | 9 +- 19 files changed, 172 insertions(+), 260 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java index 61827449e881c..4ce503298e1d0 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java +++ b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java @@ -261,29 +261,23 @@ void testAdminClient() throws Throwable { nodes.size(), testConfig.numClusterNodes); } tryFeature("createTopics", testConfig.createTopicsSupported, - new Invoker() { - @Override - public void invoke() throws Throwable { - try { - client.createTopics(Collections.singleton( - new NewTopic("newtopic", 1, (short) 1))).all().get(); - } catch (ExecutionException e) { - throw e.getCause(); - } + () -> { + try { + client.createTopics(Collections.singleton( + new NewTopic("newtopic", 1, (short) 1))).all().get(); + } catch (ExecutionException e) { + throw e.getCause(); } }, - new ResultTester() { - @Override - public void test() throws Throwable { - while (true) { - try { - client.describeTopics(Collections.singleton("newtopic")).all().get(); - break; - } catch (ExecutionException e) { - if (e.getCause() instanceof UnknownTopicOrPartitionException) - continue; - throw e; - } + () -> { + while (true) { + try { + client.describeTopics(Collections.singleton("newtopic")).all().get(); + break; + } catch (ExecutionException e) { + if (e.getCause() instanceof UnknownTopicOrPartitionException) + continue; + throw e; } } }); @@ -305,16 +299,13 @@ public void test() throws Throwable { log.info("Did not see newtopic. Retrying listTopics..."); } tryFeature("describeAclsSupported", testConfig.describeAclsSupported, - new Invoker() { - @Override - public void invoke() throws Throwable { - try { - client.describeAcls(AclBindingFilter.ANY).values().get(); - } catch (ExecutionException e) { - if (e.getCause() instanceof SecurityDisabledException) - return; - throw e.getCause(); - } + () -> { + try { + client.describeAcls(AclBindingFilter.ANY).values().get(); + } catch (ExecutionException e) { + if (e.getCause() instanceof SecurityDisabledException) + return; + throw e.getCause(); } }); } @@ -384,18 +375,8 @@ public void testConsume(final long prodTimeMs) throws Throwable { } final OffsetsForTime offsetsForTime = new OffsetsForTime(); tryFeature("offsetsForTimes", testConfig.offsetsForTimesSupported, - new Invoker() { - @Override - public void invoke() { - offsetsForTime.result = consumer.offsetsForTimes(timestampsToSearch); - } - }, - new ResultTester() { - @Override - public void test() { - log.info("offsetsForTime = {}", offsetsForTime.result); - } - }); + () -> offsetsForTime.result = consumer.offsetsForTimes(timestampsToSearch), + () -> log.info("offsetsForTime = {}", offsetsForTime.result)); // Whether or not offsetsForTimes works, beginningOffsets and endOffsets // should work. consumer.beginningOffsets(timestampsToSearch.keySet()); @@ -486,11 +467,7 @@ private interface ResultTester { } private void tryFeature(String featureName, boolean supported, Invoker invoker) throws Throwable { - tryFeature(featureName, supported, invoker, new ResultTester() { - @Override - public void test() { - } - }); + tryFeature(featureName, supported, invoker, () -> {}); } private void tryFeature(String featureName, boolean supported, Invoker invoker, ResultTester resultTester) diff --git a/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java b/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java index 0e5d1300cea63..14805ecbcf770 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java +++ b/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java @@ -32,12 +32,7 @@ public class ToolsUtils { public static void printMetrics(Map metrics) { if (metrics != null && !metrics.isEmpty()) { int maxLengthOfDisplayName = 0; - TreeMap sortedMetrics = new TreeMap<>(new Comparator() { - @Override - public int compare(String o1, String o2) { - return o1.compareTo(o2); - } - }); + TreeMap sortedMetrics = new TreeMap<>(); for (Metric metric : metrics.values()) { MetricName mName = metric.metricName(); String mergedName = mName.group() + ":" + mName.name() + ":" + mName.tags(); diff --git a/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java b/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java index 27e7c7fda5d8e..a0ac1f188f21b 100644 --- a/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java +++ b/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java @@ -263,18 +263,15 @@ public static void main(String[] args) throws IOException { final AtomicBoolean isShuttingDown = new AtomicBoolean(false); final AtomicLong remainingMessages = new AtomicLong(maxMessages); final AtomicLong numMessagesProcessed = new AtomicLong(0); - Runtime.getRuntime().addShutdownHook(new Thread() { - @Override - public void run() { - isShuttingDown.set(true); - // Flush any remaining messages - producer.close(); - synchronized (consumer) { - consumer.close(); - } - System.out.println(shutDownString(numMessagesProcessed.get(), remainingMessages.get(), transactionalId)); + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + isShuttingDown.set(true); + // Flush any remaining messages + producer.close(); + synchronized (consumer) { + consumer.close(); } - }); + System.out.println(shutDownString(numMessagesProcessed.get(), remainingMessages.get(), transactionalId)); + })); try { Random random = new Random(); diff --git a/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java b/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java index 58f34718b8afd..129784185eace 100644 --- a/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java +++ b/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java @@ -620,12 +620,7 @@ public static void main(String[] args) { try { final VerifiableConsumer consumer = createFromArgs(parser, args); - Runtime.getRuntime().addShutdownHook(new Thread() { - @Override - public void run() { - consumer.close(); - } - }); + Runtime.getRuntime().addShutdownHook(new Thread(() -> consumer.close())); consumer.run(); } catch (ArgumentParserException e) { diff --git a/tools/src/main/java/org/apache/kafka/tools/VerifiableLog4jAppender.java b/tools/src/main/java/org/apache/kafka/tools/VerifiableLog4jAppender.java index 9d23bf311768c..12aa4f45fde49 100644 --- a/tools/src/main/java/org/apache/kafka/tools/VerifiableLog4jAppender.java +++ b/tools/src/main/java/org/apache/kafka/tools/VerifiableLog4jAppender.java @@ -241,13 +241,10 @@ public static void main(String[] args) throws IOException { final VerifiableLog4jAppender appender = createFromArgs(args); boolean infinite = appender.maxMessages < 0; - Runtime.getRuntime().addShutdownHook(new Thread() { - @Override - public void run() { - // Trigger main thread to stop producing messages - appender.stopLogging = true; - } - }); + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + // Trigger main thread to stop producing messages + appender.stopLogging = true; + })); long maxMessages = infinite ? Long.MAX_VALUE : appender.maxMessages; for (long i = 0; i < maxMessages; i++) { diff --git a/tools/src/main/java/org/apache/kafka/tools/VerifiableProducer.java b/tools/src/main/java/org/apache/kafka/tools/VerifiableProducer.java index f0a991fadae88..3e6f3f14ce7a8 100644 --- a/tools/src/main/java/org/apache/kafka/tools/VerifiableProducer.java +++ b/tools/src/main/java/org/apache/kafka/tools/VerifiableProducer.java @@ -517,22 +517,19 @@ public static void main(String[] args) { final long startMs = System.currentTimeMillis(); ThroughputThrottler throttler = new ThroughputThrottler(producer.throughput, startMs); - Runtime.getRuntime().addShutdownHook(new Thread() { - @Override - public void run() { - // Trigger main thread to stop producing messages - producer.stopProducing = true; + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + // Trigger main thread to stop producing messages + producer.stopProducing = true; - // Flush any remaining messages - producer.close(); + // Flush any remaining messages + producer.close(); - // Print a summary - long stopMs = System.currentTimeMillis(); - double avgThroughput = 1000 * ((producer.numAcked) / (double) (stopMs - startMs)); + // Print a summary + long stopMs = System.currentTimeMillis(); + double avgThroughput = 1000 * ((producer.numAcked) / (double) (stopMs - startMs)); - producer.printJson(new ToolData(producer.numSent, producer.numAcked, producer.throughput, avgThroughput)); - } - }); + producer.printJson(new ToolData(producer.numSent, producer.numAcked, producer.throughput, avgThroughput)); + })); producer.run(throttler); } catch (ArgumentParserException e) { diff --git a/tools/src/main/java/org/apache/kafka/trogdor/agent/Agent.java b/tools/src/main/java/org/apache/kafka/trogdor/agent/Agent.java index 20d34b7f43cb8..c76ef263c2233 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/agent/Agent.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/agent/Agent.java @@ -147,18 +147,15 @@ public static void main(String[] args) throws Exception { log.info("Starting agent process."); final Agent agent = new Agent(platform, Scheduler.SYSTEM, restServer, resource); restServer.start(resource); - Runtime.getRuntime().addShutdownHook(new Thread() { - @Override - public void run() { - log.warn("Running agent shutdown hook."); - try { - agent.beginShutdown(); - agent.waitForShutdown(); - } catch (Exception e) { - log.error("Got exception while running agent shutdown hook.", e); - } + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + log.warn("Running agent shutdown hook."); + try { + agent.beginShutdown(); + agent.waitForShutdown(); + } catch (Exception e) { + log.error("Got exception while running agent shutdown hook.", e); } - }); + })); agent.waitForShutdown(); } }; diff --git a/tools/src/main/java/org/apache/kafka/trogdor/agent/WorkerManager.java b/tools/src/main/java/org/apache/kafka/trogdor/agent/WorkerManager.java index 59d34c90ab6fb..ef02716921247 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/agent/WorkerManager.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/agent/WorkerManager.java @@ -317,21 +317,18 @@ public void createWorker(long workerId, String taskId, TaskSpec spec) throws Thr return; } KafkaFutureImpl haltFuture = new KafkaFutureImpl<>(); - haltFuture.thenApply(new KafkaFuture.BaseFunction() { - @Override - public Void apply(String errorString) { - if (errorString == null) - errorString = ""; - if (errorString.isEmpty()) { - log.info("{}: Worker {} is halting.", nodeName, worker); - } else { - log.info("{}: Worker {} is halting with error {}", - nodeName, worker, errorString); - } - stateChangeExecutor.submit( - new HandleWorkerHalting(worker, errorString, false)); - return null; + haltFuture.thenApply((KafkaFuture.BaseFunction) errorString -> { + if (errorString == null) + errorString = ""; + if (errorString.isEmpty()) { + log.info("{}: Worker {} is halting.", nodeName, worker); + } else { + log.info("{}: Worker {} is halting with error {}", + nodeName, worker, errorString); } + stateChangeExecutor.submit( + new HandleWorkerHalting(worker, errorString, false)); + return null; }); try { worker.taskWorker.start(platform, worker.status, haltFuture); diff --git a/tools/src/main/java/org/apache/kafka/trogdor/coordinator/Coordinator.java b/tools/src/main/java/org/apache/kafka/trogdor/coordinator/Coordinator.java index cd3da904c2bfd..a41a6f2bca79f 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/coordinator/Coordinator.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/coordinator/Coordinator.java @@ -162,18 +162,15 @@ public static void main(String[] args) throws Exception { final Coordinator coordinator = new Coordinator(platform, Scheduler.SYSTEM, restServer, resource, ThreadLocalRandom.current().nextLong(0, Long.MAX_VALUE / 2)); restServer.start(resource); - Runtime.getRuntime().addShutdownHook(new Thread() { - @Override - public void run() { - log.warn("Running coordinator shutdown hook."); - try { - coordinator.beginShutdown(false); - coordinator.waitForShutdown(); - } catch (Exception e) { - log.error("Got exception while running coordinator shutdown hook.", e); - } + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + log.warn("Running coordinator shutdown hook."); + try { + coordinator.beginShutdown(false); + coordinator.waitForShutdown(); + } catch (Exception e) { + log.error("Got exception while running coordinator shutdown hook.", e); } - }); + })); coordinator.waitForShutdown(); } }; diff --git a/tools/src/main/java/org/apache/kafka/trogdor/rest/JsonRestServer.java b/tools/src/main/java/org/apache/kafka/trogdor/rest/JsonRestServer.java index ee8643b26f94d..196ec82ea8857 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/rest/JsonRestServer.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/rest/JsonRestServer.java @@ -132,22 +132,19 @@ public int port() { */ public void beginShutdown() { if (!shutdownExecutor.isShutdown()) { - shutdownExecutor.submit(new Callable() { - @Override - public Void call() throws Exception { - try { - log.info("Stopping REST server"); - jettyServer.stop(); - jettyServer.join(); - log.info("REST server stopped"); - } catch (Exception e) { - log.error("Unable to stop REST server", e); - } finally { - jettyServer.destroy(); - } - shutdownExecutor.shutdown(); - return null; + shutdownExecutor.submit((Callable) () -> { + try { + log.info("Stopping REST server"); + jettyServer.stop(); + jettyServer.join(); + log.info("REST server stopped"); + } catch (Exception e) { + log.error("Unable to stop REST server", e); + } finally { + jettyServer.destroy(); } + shutdownExecutor.shutdown(); + return null; }); } } diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java index c22396f7a22d8..e41b2b0006f36 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java @@ -98,12 +98,7 @@ public ConnectionStressAction action() { } public TaskController newController(String id) { - return new TaskController() { - @Override - public Set targetNodes(Topology topology) { - return new TreeSet<>(clientNodes); - } - }; + return topology -> new TreeSet<>(clientNodes); } @Override diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java index d15172f81e1f0..2a7bd422a5e14 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java @@ -170,12 +170,7 @@ public TopicsSpec inactiveTopics() { @Override public TaskController newController(String id) { - return new TaskController() { - @Override - public Set targetNodes(Topology topology) { - return Collections.singleton(producerNode); - } - }; + return topology -> Collections.singleton(producerNode); } @Override diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java index 669fafcc75ed5..f823256f233dc 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java @@ -248,16 +248,13 @@ public void run() { ProducerRecord record = new ProducerRecord<>(partition.topic(), partition.partition(), KEY_GENERATOR.generate(messageIndex), spec.valueGenerator().generate(messageIndex)); - producer.send(record, new Callback() { - @Override - public void onCompletion(RecordMetadata metadata, Exception exception) { - if (exception == null) { - unackedSends.countDown(); - } else { - log.info("{}: Got exception when sending message {}: {}", - id, messageIndex, exception.getMessage()); - toSendTracker.addFailed(messageIndex); - } + producer.send(record, (metadata, exception) -> { + if (exception == null) { + unackedSends.countDown(); + } else { + log.info("{}: Got exception when sending message {}: {}", + id, messageIndex, exception.getMessage()); + toSendTracker.addFailed(messageIndex); } }); } diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java index 9522e0a938e07..b63fae485da8f 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java @@ -124,12 +124,7 @@ public Map consumerConf() { @Override public TaskController newController(String id) { - return new TaskController() { - @Override - public Set targetNodes(Topology topology) { - return Collections.singleton(clientNode); - } - }; + return topology -> Collections.singleton(clientNode); } @Override diff --git a/tools/src/test/java/org/apache/kafka/trogdor/agent/AgentTest.java b/tools/src/test/java/org/apache/kafka/trogdor/agent/AgentTest.java index 158e690da4ba2..f0ea47535818f 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/agent/AgentTest.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/agent/AgentTest.java @@ -304,7 +304,7 @@ public void testKiboshFaults() throws Exception { try (MockKibosh mockKibosh = new MockKibosh()) { Assert.assertEquals(KiboshControlFile.EMPTY, mockKibosh.read()); FilesUnreadableFaultSpec fooSpec = new FilesUnreadableFaultSpec(0, 900000, - Collections.singleton("myAgent"), mockKibosh.tempDir.getPath().toString(), "/foo", 123); + Collections.singleton("myAgent"), mockKibosh.tempDir.getPath(), "/foo", 123); client.createWorker(new CreateWorkerRequest(0, "foo", fooSpec)); new ExpectedTasks(). addTask(new ExpectedTaskBuilder("foo"). @@ -314,7 +314,7 @@ public void testKiboshFaults() throws Exception { Assert.assertEquals(new KiboshControlFile(Collections.singletonList( new KiboshFilesUnreadableFaultSpec("/foo", 123))), mockKibosh.read()); FilesUnreadableFaultSpec barSpec = new FilesUnreadableFaultSpec(0, 900000, - Collections.singleton("myAgent"), mockKibosh.tempDir.getPath().toString(), "/bar", 456); + Collections.singleton("myAgent"), mockKibosh.tempDir.getPath(), "/bar", 456); client.createWorker(new CreateWorkerRequest(1, "bar", barSpec)); new ExpectedTasks(). addTask(new ExpectedTaskBuilder("foo"). diff --git a/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java b/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java index c092c92d6d602..41e8f56373a80 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java @@ -142,71 +142,65 @@ public ExpectedTasks addTask(ExpectedTask task) { } public ExpectedTasks waitFor(final CoordinatorClient client) throws InterruptedException { - TestUtils.waitForCondition(new TestCondition() { - @Override - public boolean conditionMet() { - TasksResponse tasks = null; - try { - tasks = client.tasks(new TasksRequest(null, 0, 0, 0, 0, Optional.empty())); - } catch (Exception e) { - log.info("Unable to get coordinator tasks", e); - throw new RuntimeException(e); - } - StringBuilder errors = new StringBuilder(); - for (Map.Entry entry : expected.entrySet()) { - String id = entry.getKey(); - ExpectedTask task = entry.getValue(); - String differences = task.compare(tasks.tasks().get(id)); - if (differences != null) { - errors.append(differences); - } - } - String errorString = errors.toString(); - if (!errorString.isEmpty()) { - log.info("EXPECTED TASKS: {}", JsonUtil.toJsonString(expected)); - log.info("ACTUAL TASKS : {}", JsonUtil.toJsonString(tasks.tasks())); - log.info(errorString); - return false; + TestUtils.waitForCondition(() -> { + TasksResponse tasks = null; + try { + tasks = client.tasks(new TasksRequest(null, 0, 0, 0, 0, Optional.empty())); + } catch (Exception e) { + log.info("Unable to get coordinator tasks", e); + throw new RuntimeException(e); + } + StringBuilder errors = new StringBuilder(); + for (Map.Entry entry : expected.entrySet()) { + String id = entry.getKey(); + ExpectedTask task = entry.getValue(); + String differences = task.compare(tasks.tasks().get(id)); + if (differences != null) { + errors.append(differences); } - return true; } + String errorString = errors.toString(); + if (!errorString.isEmpty()) { + log.info("EXPECTED TASKS: {}", JsonUtil.toJsonString(expected)); + log.info("ACTUAL TASKS : {}", JsonUtil.toJsonString(tasks.tasks())); + log.info(errorString); + return false; + } + return true; }, "Timed out waiting for expected tasks " + JsonUtil.toJsonString(expected)); return this; } public ExpectedTasks waitFor(final AgentClient client) throws InterruptedException { - TestUtils.waitForCondition(new TestCondition() { - @Override - public boolean conditionMet() { - AgentStatusResponse status = null; - try { - status = client.status(); - } catch (Exception e) { - log.info("Unable to get agent status", e); - throw new RuntimeException(e); - } - StringBuilder errors = new StringBuilder(); - HashMap taskIdToWorkerState = new HashMap<>(); - for (WorkerState state : status.workers().values()) { - taskIdToWorkerState.put(state.taskId(), state); - } - for (Map.Entry entry : expected.entrySet()) { - String id = entry.getKey(); - ExpectedTask worker = entry.getValue(); - String differences = worker.compare(taskIdToWorkerState.get(id)); - if (differences != null) { - errors.append(differences); - } - } - String errorString = errors.toString(); - if (!errorString.isEmpty()) { - log.info("EXPECTED WORKERS: {}", JsonUtil.toJsonString(expected)); - log.info("ACTUAL WORKERS : {}", JsonUtil.toJsonString(status.workers())); - log.info(errorString); - return false; + TestUtils.waitForCondition(() -> { + AgentStatusResponse status = null; + try { + status = client.status(); + } catch (Exception e) { + log.info("Unable to get agent status", e); + throw new RuntimeException(e); + } + StringBuilder errors = new StringBuilder(); + HashMap taskIdToWorkerState = new HashMap<>(); + for (WorkerState state : status.workers().values()) { + taskIdToWorkerState.put(state.taskId(), state); + } + for (Map.Entry entry : expected.entrySet()) { + String id = entry.getKey(); + ExpectedTask worker = entry.getValue(); + String differences = worker.compare(taskIdToWorkerState.get(id)); + if (differences != null) { + errors.append(differences); } - return true; } + String errorString = errors.toString(); + if (!errorString.isEmpty()) { + log.info("EXPECTED WORKERS: {}", JsonUtil.toJsonString(expected)); + log.info("ACTUAL WORKERS : {}", JsonUtil.toJsonString(status.workers())); + log.info(errorString); + return false; + } + return true; }, "Timed out waiting for expected workers " + JsonUtil.toJsonString(expected)); return this; } diff --git a/tools/src/test/java/org/apache/kafka/trogdor/common/MiniTrogdorCluster.java b/tools/src/test/java/org/apache/kafka/trogdor/common/MiniTrogdorCluster.java index 46315c27d15db..9edffaa757d27 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/common/MiniTrogdorCluster.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/common/MiniTrogdorCluster.java @@ -172,27 +172,24 @@ public MiniTrogdorCluster build() throws Exception { ThreadUtils.createThreadFactory("MiniTrogdorClusterStartupThread%d", false)); final AtomicReference failure = new AtomicReference(null); for (final Map.Entry entry : nodes.entrySet()) { - executor.submit(new Callable() { - @Override - public Void call() throws Exception { - String nodeName = entry.getKey(); - try { - NodeData node = entry.getValue(); - node.platform = new BasicPlatform(nodeName, topology, scheduler, commandRunner); - if (node.agentRestResource != null) { - node.agent = new Agent(node.platform, scheduler, node.agentRestServer, - node.agentRestResource); - } - if (node.coordinatorRestResource != null) { - node.coordinator = new Coordinator(node.platform, scheduler, - node.coordinatorRestServer, node.coordinatorRestResource, 0); - } - } catch (Exception e) { - log.error("Unable to initialize {}", nodeName, e); - failure.compareAndSet(null, e); + executor.submit((Callable) () -> { + String nodeName = entry.getKey(); + try { + NodeData node = entry.getValue(); + node.platform = new BasicPlatform(nodeName, topology, scheduler, commandRunner); + if (node.agentRestResource != null) { + node.agent = new Agent(node.platform, scheduler, node.agentRestServer, + node.agentRestResource); } - return null; + if (node.coordinatorRestResource != null) { + node.coordinator = new Coordinator(node.platform, scheduler, + node.coordinatorRestServer, node.coordinatorRestResource, 0); + } + } catch (Exception e) { + log.error("Unable to initialize {}", nodeName, e); + failure.compareAndSet(null, e); } + return null; }); } executor.shutdown(); diff --git a/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java b/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java index 0207104458dcf..c2f927b2f1d6f 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java @@ -318,12 +318,8 @@ public ExpectedLines addLine(String line) { public ExpectedLines waitFor(final String nodeName, final CapturingCommandRunner runner) throws InterruptedException { - TestUtils.waitForCondition(new TestCondition() { - @Override - public boolean conditionMet() { - return linesMatch(nodeName, runner.lines(nodeName)); - } - }, "failed to find the expected lines " + this.toString()); + TestUtils.waitForCondition(() -> linesMatch(nodeName, runner.lines(nodeName)), + "failed to find the expected lines " + this.toString()); return this; } @@ -473,7 +469,7 @@ public void testTasksRequest() throws Exception { assertEquals(0, coordinatorClient.tasks( new TasksRequest(null, 10, 0, 10, 0, Optional.empty())).tasks().size()); TasksResponse resp1 = coordinatorClient.tasks( - new TasksRequest(Arrays.asList(new String[] {"foo", "baz" }), 0, 0, 0, 0, Optional.empty())); + new TasksRequest(Arrays.asList("foo", "baz"), 0, 0, 0, 0, Optional.empty())); assertTrue(resp1.tasks().containsKey("foo")); assertFalse(resp1.tasks().containsKey("bar")); assertEquals(1, resp1.tasks().size()); diff --git a/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java b/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java index ade055d3d0683..cdfd3e3243e14 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java @@ -53,12 +53,9 @@ public synchronized void start(Platform platform, WorkerStatusTracker status, if (exitMs == null) { exitMs = Long.MAX_VALUE; } - this.future = platform.scheduler().schedule(executor, new Callable() { - @Override - public Void call() throws Exception { - haltFuture.complete(spec.error()); - return null; - } + this.future = platform.scheduler().schedule(executor, () -> { + haltFuture.complete(spec.error()); + return null; }, exitMs); } From c8fbd86677cd59bcf7f7883451adb85d19f165b0 Mon Sep 17 00:00:00 2001 From: Srinivas Reddy Date: Wed, 19 Dec 2018 21:17:57 +0800 Subject: [PATCH 2/4] Fix style check errors --- .../kafka/tools/ClientCompatibilityTest.java | 65 ++++++++++++------- .../org/apache/kafka/tools/ToolsUtils.java | 1 - .../workload/ConnectionStressSpec.java | 2 - .../trogdor/workload/ProduceBenchSpec.java | 2 - .../trogdor/workload/RoundTripWorker.java | 2 - .../workload/RoundTripWorkloadSpec.java | 2 - .../kafka/trogdor/common/ExpectedTasks.java | 1 - .../trogdor/coordinator/CoordinatorTest.java | 1 - .../kafka/trogdor/task/SampleTaskWorker.java | 1 - 9 files changed, 43 insertions(+), 34 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java index 4ce503298e1d0..1f50f62496d9c 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java +++ b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java @@ -269,35 +269,21 @@ void testAdminClient() throws Throwable { throw e.getCause(); } }, - () -> { - while (true) { - try { - client.describeTopics(Collections.singleton("newtopic")).all().get(); - break; - } catch (ExecutionException e) { - if (e.getCause() instanceof UnknownTopicOrPartitionException) - continue; - throw e; - } - } - }); + () -> new CreateTopicsResultTester(client, Collections.singleton("newtopic")).invoke() + ); + while (true) { Collection listings = client.listTopics().listings().get(); if (!testConfig.createTopicsSupported) break; - boolean foundNewTopic = false; - for (TopicListing listing : listings) { - if (listing.name().equals("newtopic")) { - if (listing.isInternal()) - throw new KafkaException("Did not expect newtopic to be an internal topic."); - foundNewTopic = true; - } - } - if (foundNewTopic) + + if (isFoundNewTopic(listings, "newtopic")) break; + Thread.sleep(1); log.info("Did not see newtopic. Retrying listTopics..."); } + tryFeature("describeAclsSupported", testConfig.describeAclsSupported, () -> { try { @@ -311,6 +297,18 @@ void testAdminClient() throws Throwable { } } + private boolean isFoundNewTopic(Collection listings, String newTopicName) { + boolean foundNewTopic = false; + for (TopicListing listing : listings) { + if (listing.name().equals(newTopicName)) { + if (listing.isInternal()) + throw new KafkaException("Did not expect newtopic to be an internal topic."); + foundNewTopic = true; + } + } + return foundNewTopic; + } + private static class OffsetsForTime { Map result; @@ -467,7 +465,7 @@ private interface ResultTester { } private void tryFeature(String featureName, boolean supported, Invoker invoker) throws Throwable { - tryFeature(featureName, supported, invoker, () -> {}); + tryFeature(featureName, supported, invoker, () -> { }); } private void tryFeature(String featureName, boolean supported, Invoker invoker, ResultTester resultTester) @@ -487,4 +485,27 @@ private void tryFeature(String featureName, boolean supported, Invoker invoker, } resultTester.test(); } + + private class CreateTopicsResultTester { + private final Collection topics; + private AdminClient client; + + public CreateTopicsResultTester(AdminClient client, Collection topics) { + this.client = client; + this.topics = topics; + } + + public void invoke() throws InterruptedException, ExecutionException { + while (true) { + try { + client.describeTopics(topics).all().get(); + break; + } catch (ExecutionException e) { + if (e.getCause() instanceof UnknownTopicOrPartitionException) + continue; + throw e; + } + } + } + } } diff --git a/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java b/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java index 14805ecbcf770..3a80b5811f3fd 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java +++ b/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java @@ -19,7 +19,6 @@ import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; -import java.util.Comparator; import java.util.Map; import java.util.TreeMap; diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java index e41b2b0006f36..6141d30b60b6b 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressSpec.java @@ -19,7 +19,6 @@ import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonProperty; -import org.apache.kafka.trogdor.common.Topology; import org.apache.kafka.trogdor.task.TaskController; import org.apache.kafka.trogdor.task.TaskSpec; import org.apache.kafka.trogdor.task.TaskWorker; @@ -28,7 +27,6 @@ import java.util.Collections; import java.util.List; import java.util.Map; -import java.util.Set; import java.util.TreeSet; /** diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java index 2a7bd422a5e14..34b5393f05858 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchSpec.java @@ -19,7 +19,6 @@ import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonProperty; -import org.apache.kafka.trogdor.common.Topology; import org.apache.kafka.trogdor.task.TaskController; import org.apache.kafka.trogdor.task.TaskSpec; import org.apache.kafka.trogdor.task.TaskWorker; @@ -27,7 +26,6 @@ import java.util.Collections; import java.util.Map; import java.util.Optional; -import java.util.Set; /** * The specification for a benchmark that produces messages to a set of topics. diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java index f823256f233dc..b22292adf7f9f 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java @@ -25,11 +25,9 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; -import org.apache.kafka.clients.producer.Callback; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.errors.TimeoutException; diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java index b63fae485da8f..42e09ee798f3a 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorkloadSpec.java @@ -19,14 +19,12 @@ import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonProperty; -import org.apache.kafka.trogdor.common.Topology; import org.apache.kafka.trogdor.task.TaskController; import org.apache.kafka.trogdor.task.TaskSpec; import org.apache.kafka.trogdor.task.TaskWorker; import java.util.Collections; import java.util.Map; -import java.util.Set; /** * The specification for a workload that sends messages to a broker and then diff --git a/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java b/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java index 41e8f56373a80..3eb781c4f7fdb 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/common/ExpectedTasks.java @@ -19,7 +19,6 @@ import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonProperty; -import org.apache.kafka.test.TestCondition; import org.apache.kafka.test.TestUtils; import org.apache.kafka.trogdor.agent.AgentClient; import org.apache.kafka.trogdor.coordinator.CoordinatorClient; diff --git a/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java b/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java index c2f927b2f1d6f..db1afac3ea484 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/coordinator/CoordinatorTest.java @@ -24,7 +24,6 @@ import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Scheduler; import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.test.TestCondition; import org.apache.kafka.test.TestUtils; import org.apache.kafka.trogdor.agent.AgentClient; import org.apache.kafka.trogdor.common.CapturingCommandRunner; diff --git a/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java b/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java index cdfd3e3243e14..404817a708982 100644 --- a/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java +++ b/tools/src/test/java/org/apache/kafka/trogdor/task/SampleTaskWorker.java @@ -22,7 +22,6 @@ import org.apache.kafka.trogdor.common.Platform; import org.apache.kafka.trogdor.common.ThreadUtils; -import java.util.concurrent.Callable; import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; From 245af8bbf822f0b8a572e4c73e3b91e7fee08e2d Mon Sep 17 00:00:00 2001 From: Srinivas Reddy Date: Thu, 20 Dec 2018 00:07:22 +0800 Subject: [PATCH 3/4] static inner class to improve memory performance --- .../java/org/apache/kafka/tools/ClientCompatibilityTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java index 1f50f62496d9c..47682206961ed 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java +++ b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java @@ -486,7 +486,7 @@ private void tryFeature(String featureName, boolean supported, Invoker invoker, resultTester.test(); } - private class CreateTopicsResultTester { + private static class CreateTopicsResultTester { private final Collection topics; private AdminClient client; From 02c467b16fa69527c3ae990aa56eaa6175b168d5 Mon Sep 17 00:00:00 2001 From: Srinivas Reddy Date: Thu, 20 Dec 2018 19:32:47 +0800 Subject: [PATCH 4/4] changed the class to method --- .../kafka/tools/ClientCompatibilityTest.java | 53 ++++++++----------- 1 file changed, 22 insertions(+), 31 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java index 47682206961ed..5b7e2287e95fe 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java +++ b/tools/src/main/java/org/apache/kafka/tools/ClientCompatibilityTest.java @@ -269,7 +269,7 @@ void testAdminClient() throws Throwable { throw e.getCause(); } }, - () -> new CreateTopicsResultTester(client, Collections.singleton("newtopic")).invoke() + () -> createTopicsResultTest(client, Collections.singleton("newtopic")) ); while (true) { @@ -277,7 +277,7 @@ void testAdminClient() throws Throwable { if (!testConfig.createTopicsSupported) break; - if (isFoundNewTopic(listings, "newtopic")) + if (topicExists(listings, "newtopic")) break; Thread.sleep(1); @@ -297,16 +297,30 @@ void testAdminClient() throws Throwable { } } - private boolean isFoundNewTopic(Collection listings, String newTopicName) { - boolean foundNewTopic = false; + private void createTopicsResultTest(AdminClient client, Collection topics) + throws InterruptedException, ExecutionException { + while (true) { + try { + client.describeTopics(topics).all().get(); + break; + } catch (ExecutionException e) { + if (e.getCause() instanceof UnknownTopicOrPartitionException) + continue; + throw e; + } + } + } + + private boolean topicExists(Collection listings, String topicName) { + boolean foundTopic = false; for (TopicListing listing : listings) { - if (listing.name().equals(newTopicName)) { + if (listing.name().equals(topicName)) { if (listing.isInternal()) - throw new KafkaException("Did not expect newtopic to be an internal topic."); - foundNewTopic = true; + throw new KafkaException(String.format("Did not expect %s to be an internal topic.", topicName)); + foundTopic = true; } } - return foundNewTopic; + return foundTopic; } private static class OffsetsForTime { @@ -485,27 +499,4 @@ private void tryFeature(String featureName, boolean supported, Invoker invoker, } resultTester.test(); } - - private static class CreateTopicsResultTester { - private final Collection topics; - private AdminClient client; - - public CreateTopicsResultTester(AdminClient client, Collection topics) { - this.client = client; - this.topics = topics; - } - - public void invoke() throws InterruptedException, ExecutionException { - while (true) { - try { - client.describeTopics(topics).all().get(); - break; - } catch (ExecutionException e) { - if (e.getCause() instanceof UnknownTopicOrPartitionException) - continue; - throw e; - } - } - } - } }