Fehlerbehandlung in EndlessConsumer.destroy() korrigiert
[demos/kafka/training] / src / main / java / de / juplo / kafka / EndlessConsumer.java
index da2f8f0..f25e93c 100644 (file)
@@ -25,6 +25,7 @@ public class EndlessConsumer implements Runnable
   private final String groupId;
   private final String id;
   private final String topic;
+  private final String autoOffsetReset;
 
   private AtomicBoolean running = new AtomicBoolean();
   private long consumed = 0;
@@ -36,13 +37,15 @@ public class EndlessConsumer implements Runnable
       String bootstrapServer,
       String groupId,
       String clientId,
-      String topic)
+      String topic,
+      String autoOffsetReset)
   {
     this.executor = executor;
     this.bootstrapServer = bootstrapServer;
     this.groupId = groupId;
     this.id = clientId;
     this.topic = topic;
+    this.autoOffsetReset = autoOffsetReset;
   }
 
   @Override
@@ -54,7 +57,7 @@ public class EndlessConsumer implements Runnable
       props.put("bootstrap.servers", bootstrapServer);
       props.put("group.id", groupId);
       props.put("client.id", id);
-      props.put("auto.offset.reset", "earliest");
+      props.put("auto.offset.reset", autoOffsetReset);
       props.put("key.deserializer", StringDeserializer.class.getName());
       props.put("value.deserializer", StringDeserializer.class.getName());
 
@@ -107,7 +110,7 @@ public class EndlessConsumer implements Runnable
   {
     boolean stateChanged = running.compareAndSet(false, true);
     if (!stateChanged)
-      throw new RuntimeException("Consumer instance " + id + " is already running!");
+      throw new IllegalStateException("Consumer instance " + id + " is already running!");
 
     log.info("{} - Starting - consumed {} messages before", id, consumed);
     future = executor.submit(this);
@@ -117,7 +120,7 @@ public class EndlessConsumer implements Runnable
   {
     boolean stateChanged = running.compareAndSet(true, false);
     if (!stateChanged)
-      throw new RuntimeException("Consumer instance " + id + " is not running!");
+      throw new IllegalStateException("Consumer instance " + id + " is not running!");
 
     log.info("{} - Stopping", id);
     consumer.wakeup();
@@ -137,6 +140,10 @@ public class EndlessConsumer implements Runnable
     {
       log.info("{} - Was already stopped", id);
     }
+    catch (Exception e)
+    {
+      log.error("{} - Unexpected exception while trying to stop the consumer", id, e);
+    }
     finally
     {
       log.info("{}: Consumed {} messages in total, exiting!", id, consumed);