diff --git a/flink-contrib/flink-storm/src/test/java/org/apache/flink/storm/wrappers/WrapperSetupHelperTest.java b/flink-contrib/flink-storm/src/test/java/org/apache/flink/storm/wrappers/WrapperSetupHelperTest.java index ab846aff04697..eac2999d47aca 100644 --- a/flink-contrib/flink-storm/src/test/java/org/apache/flink/storm/wrappers/WrapperSetupHelperTest.java +++ b/flink-contrib/flink-storm/src/test/java/org/apache/flink/storm/wrappers/WrapperSetupHelperTest.java @@ -182,20 +182,15 @@ public void testCreateTopologyContext() { .shuffleGrouping("bolt2", TestDummyBolt.groupingStreamId) .shuffleGrouping("bolt2", TestDummyBolt.shuffleStreamId); - int counter = 0; - while (true) { - LocalCluster cluster = new LocalCluster(); - Config c = new Config(); - c.setNumAckers(0); - c.setDebug(true); - cluster.submitTopology("test", c, builder.createTopology()); - Utils.sleep(++counter * 3000); - cluster.shutdown(); - - if (TestSink.result.size() == 8) { - break; - } + LocalCluster cluster = new LocalCluster(); + Config c = new Config(); + c.setNumAckers(0); + cluster.submitTopology("test", c, builder.createTopology()); + + while (TestSink.result.size() != 8) { + Utils.sleep(1); } + cluster.shutdown(); final FlinkTopology flinkBuilder = FlinkTopology.createTopology(builder); StormTopology stormTopology = flinkBuilder.getStormTopology();