private final String id;
private final String topic;
private final Consumer<K, V> consumer;
- private final ConsumerRebalanceListener rebalanceListener;
+ private final PollIntervalAwareConsumerRebalanceListener pollIntervalAwareRebalanceListener;
private final RecordHandler<K, V> handler;
private final Lock lock = new ReentrantLock();
try
{
log.info("{} - Subscribing to topic {}", id, topic);
- consumer.subscribe(Arrays.asList(topic), rebalanceListener);
+ consumer.subscribe(Arrays.asList(topic), pollIntervalAwareRebalanceListener);
while (true)
{
consumed++;
}
- handler.beforeNextPoll();
+ pollIntervalAwareRebalanceListener.beforeNextPoll();
}
}
catch(WakeupException e)