package de.juplo.kafka;
-import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
ApplicationRecordHandler recordHandler,
AdderResults adderResults,
StateRepository stateRepository,
- Consumer<String, String> consumer,
ApplicationProperties properties)
{
return new ApplicationRebalanceListener(
recordHandler,
adderResults,
stateRepository,
- properties.getClientId(),
- consumer);
+ properties.getClientId());
}
@Bean
Properties props = new Properties();
props.put("bootstrap.servers", properties.getBootstrapServer());
- props.put("partition.assignment.strategy", "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
+ props.put("partition.assignment.strategy", "org.apache.kafka.clients.consumer.StickyAssignor");
props.put("group.id", properties.getGroupId());
props.put("client.id", properties.getClientId());
props.put("auto.offset.reset", properties.getAutoOffsetReset());