【发布时间】:2020-06-23 22:05:31
【问题描述】:
我有一个简单的要求来读取 kafka 消息并存储在数据库中。我在批处理侦听器模式下使用 spring kafka。我已经浏览了 spring kafka 文档,但仍然不清楚在批处理侦听器模式下使用 spring kafka 时,它是否以批处理模式提交 db 事务,如果失败,是否会回滚整个事务?
如果失败,它会再次寻找相同的记录集吗?
我有以下配置,
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfigProperties.getBootstrapservers());
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaConfigProperties.getConsumer().getGroupid());
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, kafkaConfigProperties.getConsumer().getOffset());
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,250);
props.put(ApplicationConstant.KAFKA_SCHEMA_URL_PROPERTY, kafkaConfigProperties.getSchemaregistry());
@Bean
public ConcurrentKafkaListenerContainerFactory<String, GenericRecord> kafkaListenerContainerFactory(KafkaConfigProperties kafkaConfigProperties) {
ConcurrentKafkaListenerContainerFactory<String, GenericRecord> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory(kafkaConfigProperties));
factory.setConcurrency(2);
factory.setBatchListener(true);
ContainerProperties containerProperties = factory.getContainerProperties();
containerProperties.setAckOnError(false);
containerProperties.setAckMode(AckMode.BATCH);
return factory;
}
【问题讨论】:
标签: java spring apache-kafka spring-kafka