try
{
log.info("{} - Subscribing to topic {}", id, topic);
- rebalanceListener.enableCommits();
consumer.subscribe(Arrays.asList(topic), rebalanceListener);
while (true)
catch(WakeupException e)
{
log.info("{} - RIIING! Request to stop consumption - commiting current offsets!", id);
+ consumer.commitSync();
shutdown();
}
catch(RecordDeserializationException e)
offset,
e.getCause().toString());
+ consumer.commitSync();
shutdown(e);
}
catch(Exception e)
{
- log.error("{} - Unexpected error: {}, disabling commits", id, e.toString(), e);
- rebalanceListener.disableCommits();
+ log.error("{} - Unexpected error: {}", id, e.toString(), e);
shutdown(e);
}
finally