props.put("bootstrap.servers", broker);
props.put("group.id", groupId); // ID für die Offset-Commits
props.put("client.id", clientId); // Nur zur Wiedererkennung
+ props.put("auto.offset.reset", "earliest"); // Von Beginn an lesen
props.put("partition.assignment.strategy", "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
}
catch(Exception e)
{
- log.error("{} - Unexpected error: {}", id, e.toString());
+ log.error("{} - Unexpected error: {}, unsubscribing!", id, e.toString());
+ consumer.unsubscribe();
}
finally
{