public class EndlessConsumer implements Runnable
{
private final ExecutorService executor;
+ private final PartitionStatisticsRepository repository;
private final String bootstrapServer;
private final String groupId;
private final String id;
private final Lock lock = new ReentrantLock();
private final Condition condition = lock.newCondition();
private boolean running = false;
+ private Exception exception;
private long consumed = 0;
private KafkaConsumer<String, String> consumer = null;
public EndlessConsumer(
ExecutorService executor,
+ PartitionStatisticsRepository repository,
String bootstrapServer,
String groupId,
String clientId,
String autoOffsetReset)
{
this.executor = executor;
+ this.repository = repository;
this.bootstrapServer = bootstrapServer;
this.groupId = groupId;
this.id = clientId;
tp.partition(),
key);
}
+ repository.save(new StatisticsDocument(tp.partition(), removed));
});
}
partitions.forEach(tp ->
{
log.info("{} - adding partition: {}", id, tp);
- seen.put(tp.partition(), new HashMap<>());
+ seen.put(
+ tp.partition(),
+ repository
+ .findById(Integer.toString(tp.partition()))
+ .map(document -> document.statistics)
+ .orElse(new HashMap<>()));
});
}
});
catch(Exception e)
{
log.error("{} - Unexpected error: {}", id, e.toString(), e);
- shutdown();
+ shutdown(e);
}
finally
{
}
private void shutdown()
+ {
+ shutdown(null);
+ }
+
+ private void shutdown(Exception e)
{
lock.lock();
try
{
running = false;
+ exception = e;
condition.signal();
}
finally
log.info("{} - Starting - consumed {} messages before", id, consumed);
running = true;
+ exception = null;
executor.submit(this);
}
finally
log.info("{}: Consumed {} messages in total, exiting!", id, consumed);
}
}
+
+ public boolean running()
+ {
+ lock.lock();
+ try
+ {
+ return running;
+ }
+ finally
+ {
+ lock.unlock();
+ }
+ }
+
+ public Optional<Exception> exitStatus()
+ {
+ lock.lock();
+ try
+ {
+ if (running)
+ throw new IllegalStateException("No exit-status available: Consumer instance " + id + " is running!");
+
+ return Optional.ofNullable(exception);
+ }
+ finally
+ {
+ lock.unlock();
+ }
+ }
}