【发布时间】:2021-12-13 16:41:34
【问题描述】:
我们的基础架构必须管理 2 个 Spring Boot 应用程序之间的信息交换。 一些信息基于 Oracle DB。 我们使用 Kafka 来通知哪些信息必须由接收者管理。
(注意:我必须概括代码)
我们有 Kafka 生产者,配置如下:
@Configuration
@Lazy(false)
public class KafkaConfig {
//...
@Bean(name="trManagerJpaKafka")
public ChainedKafkaTransactionManager<Object, Object> chainedTm(KafkaTransactionManager<String, String> ktm,
JpaTransactionManager jpaTransactionManager) {
return new ChainedKafkaTransactionManager<>(jpaTransactionManager,ktm);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(KafkaOperations<?, ?> template,
ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
ConsumerFactory<Object, Object> kafkaConsumerFactory, ChainedKafkaTransactionManager<Object, Object> chainedTm) {
ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
configurer.configure(factory, kafkaConsumerFactory);
//...
factory.setErrorHandler(new SeekToCurrentErrorHandler(recoverer,new FixedBackOff(0L, 2)));
factory.getContainerProperties().setTransactionManager(chainedTm);
return factory;
}
}
然后我们有一些服务产生 Kafka 消息,像这样
public class SomeService implements ISomeService {
@Override
@Transactional(transactionManager = "trManagerJpaKafka" , rollbackFor = { Exception.class })
public boolean createMessage(Long id, CreateRDM createRDM) {
...
Entity entity = new Entity();
entity.setSomething(something);
entity.setSomethingElse(somethingElse);
entity = entityRepository.save(entity);
JSONObject message = new JSONObject();
try {
message.put("id", entity.getId());
kafkaProducerService.sendMessage(message.toString());
...
}
...
}
public class SomeGenericKafkaService implements ISomeGenericKafkaService {
...
public void sendMessage(String data) {
Map<String, Object> headers = new HashMap<String, Object>();
headers.put(KafkaHeaders.TOPIC, someTopic);
headers.put(KafkaHeaders.MESSAGE_KEY, UUID.randomUUID().toString());
KafkaProducerCallback callback = new KafkaProducerCallback();
kafkaTemplate.send(new GenericMessage<String>(data, headers)).addCallback(callback);
if (callback.isError()) {
throw new KafkaException(callback.getThrowable());
}
}
...
}
在消费者 Spring Boot 应用程序中,我们有时会遇到这个问题: 它使用来自 Kafka 的消息,但有时 DB 上的记录尚不存在,因此当我们尝试访问 DB 记录时出现异常...... 生产者中的事务对于 Jpa 和 Kafka 都不是原子的? Kafka 消息在 Jpa 插入之前提交?
同时,作为紧急解决方案,为了管理从 Oracle 检索数据的异常,我们在接收器中放置了一些重试,类似于他的:
...
Optional<Entity> optionalEntity = entityRepository.findById(id);
if (!optionalEntity.isPresent()) {
for (int i = 0; i < kafkaDbMaxRetry; i++) {
Thread.sleep(kafkaMessageDelay);
optionalEntity = entityRepository.findById(id);
if (optionalEntity.isPresent()) break;
}
}
entity = optionalEntity.get();
...
如何调整生产者,避免消费者出现这种糟糕的代码?
【问题讨论】:
标签: java spring spring-boot spring-data-jpa spring-kafka