I have a spring-boot (2.1.3) service publishing messages to a kafka(2.12-2.3.0) topic. The service creates the topic and later, after the service is up, sets the retention.ms to 1 second.
Currently debugging this code
@SpringBootApplication()
@EnableAsync
public class MetricsMsApplication {
public static void main(String[] args) {
SpringApplication.run(MetricsMsApplication.class, args);
}
@Bean
public NewTopic topic1() {
NewTopic nt = new NewTopic("metrics", 10, (short) 1);
return nt;
}
@EventListener(ApplicationReadyEvent.class)
private void init() throws ExecutionException, InterruptedException {
Map<String, Object> config = new HashMap<>();
config.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");
AdminClient client = AdminClient.create(config);
ConfigResource resource = new ConfigResource(ConfigResource.Type.TOPIC, "metrics");
// Update the retention.ms value
ConfigEntry retentionEntry = new ConfigEntry(TopicConfig.RETENTION_MS_CONFIG, "1000");
Map<ConfigResource, Config> updateConfig = new HashMap<ConfigResource, Config>();
updateConfig.put(resource, new Config(Collections.singleton(retentionEntry)));
AlterConfigsResult alterConfigsResult = client.alterConfigs(updateConfig);
alterConfigsResult.all();
}
}
I send in a couple of messages and count to 5, then start a console consumer
kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic admst-metrics --from-beginning
and still get the messages that should have expired.
The kafka logs show the retention.ms config was applied. I added cleanup.policy and set it to delete, but that shouldn't be necessary as it's the default.
What will make these messages get deleted?