@@ -22,9 +22,9 @@ class TestDroppingStaleMessages {
2222 val config: MutableMap <String , Any > = mutableMapOf ()
2323 config[Config .TOPOLOGY_MAX_SPOUT_PENDING ] = 1L
2424 config[KumulusTopology .CONF_THREAD_POOL_CORE_SIZE ] = 5L
25- config[KumulusTopology .CONF_LATE_MESSAGES_DROPPING_SHOULD_DROP ] = true ;
26- config[KumulusTopology .CONF_LATE_MESSAGES_DROPPING_STREAMS_NAME ] = setOf (" default" );
27- config[KumulusTopology .CONF_LATE_MESSAGES_DROPPING_MAX_WAIT_SECONDS ] = 1 ;
25+ config[KumulusTopology .CONF_LATE_MESSAGES_DROPPING_SHOULD_DROP ] = true
26+ config[KumulusTopology .CONF_LATE_MESSAGES_DROPPING_STREAMS_NAME ] = setOf (" default" )
27+ config[KumulusTopology .CONF_LATE_MESSAGES_DROPPING_MAX_WAIT_SECONDS ] = 1
2828
2929 builder.setSpout(" spout" , LatencyDeltaSpout ())
3030
@@ -34,7 +34,6 @@ class TestDroppingStaleMessages {
3434 builder.setBolt(" delay-unanchored-bolt" , StuckBolt ())
3535 .noneGrouping(" unanchoring-bolt" )
3636
37-
3837 val stormTopology = builder.createTopology()!!
3938 val kumulusTopology =
4039 KumulusStormTransformer .initializeTopology(stormTopology, config, " test" )
@@ -50,8 +49,8 @@ class TestDroppingStaleMessages {
5049 Thread .sleep(5000 )
5150 kumulusTopology.stop()
5251
53- logger.info { " Dropped ${ lateHookCalled} messages" }
54- assertTrue { lateHookCalled}
52+ logger.info { " Dropped $lateHookCalled messages" }
53+ assertTrue { lateHookCalled }
5554 }
5655
5756 class LatencyDeltaSpout : DummySpout ({
@@ -110,7 +109,7 @@ class TestDroppingStaleMessages {
110109 }) {
111110 override fun execute (input : Tuple , collector : BasicOutputCollector ) {
112111 logger.info { " StuckBolt: started" }
113- while (true ){
112+ while (true ) {
114113 Thread .sleep(50 )
115114 }
116115 }
0 commit comments