@Slf4j
@RequiredArgsConstructor
-public class EndlessConsumer<K, V> implements ConsumerRebalanceListener, Runnable
+public class EndlessConsumer<K, V> implements Runnable
{
private final ExecutorService executor;
private final String id;
private final String topic;
private final Consumer<K, V> consumer;
+ private final PollIntervalAwareConsumerRebalanceListener pollIntervalAwareRebalanceListener;
private final RecordHandler<K, V> handler;
private final Lock lock = new ReentrantLock();
private long consumed = 0;
- @Override
- public void onPartitionsRevoked(Collection<TopicPartition> partitions)
- {
- partitions.forEach(tp -> handler.onPartitionRevoked(tp));
- }
-
- @Override
- public void onPartitionsAssigned(Collection<TopicPartition> partitions)
- {
- partitions.forEach(tp -> handler.onPartitionAssigned(tp));
- }
-
@Override
public void run()
try
{
log.info("{} - Subscribing to topic {}", id, topic);
- consumer.subscribe(Arrays.asList(topic), this);
+ pollIntervalAwareRebalanceListener.enableCommits();
+ consumer.subscribe(Arrays.asList(topic), pollIntervalAwareRebalanceListener);
while (true)
{
consumed++;
}
- handler.beforeNextPoll();
+ pollIntervalAwareRebalanceListener.beforeNextPoll();
}
}
catch(WakeupException e)
}
catch(Exception e)
{
- log.error("{} - Unexpected error: {}", id, e.toString(), e);
+ log.error("{} - Unexpected error: {}, disabling commits", id, e.toString(), e);
+ pollIntervalAwareRebalanceListener.disableCommits();
shutdown(e);
}
finally