【发布时间】:2021-03-23 18:00:47
【问题描述】:
我正在尝试为Kafka 处理程序编写RetryAdvice;并回退到以 RecoveryCallback 的形式保存到 MongoDB。
@Bean(name = "kafkaSuccessChannel")
public ExecutorChannel kafkaSuccessChannel() {
return MessageChannels.executor("kafkaSuccessChannel", asyncExecutor()).get();
}
@Bean(name = "kafkaErrorChannel")
public ExecutorChannel kafkaErrorChannel() {
return MessageChannels.executor("kafkaSuccessChannel", asyncExecutor()).get();
}
@Bean
@ServiceActivator(inputChannel = "kafkaPublishChannel")
public KafkaProducerMessageHandler<String, String> kafkaProducerMessageHandler(
@Autowired ExecutorChannel kafkaSuccessChannel,
@Autowired RequestHandlerRetryAdvice retryAdvice) {
KafkaProducerMessageHandler<String, String> handler = new KafkaProducerMessageHandler<>(kafkaTemplate());
handler.setHeaderMapper(mapper());
handler.setLoggingEnabled(TRUE);
handler.setTopicExpression(
new SpelExpressionParser()
.parseExpression(
"headers['" + upstreamTypeHeader + "'] + '_' + headers['" + upstreamInstanceHeader + "']"));
handler.setSendSuccessChannel(kafkaSuccessChannel);
handler.setAdviceChain(Arrays.asList(retryAdvice));
// sync true implies that this Kafka handler will wait for results of kafka operations; to be used only for testing purposes.
handler.setSync(testMode);
return handler;
}
然后我在同一个类中配置如下建议
@Bean
public RequestHandlerRetryAdvice retryAdvice(@Autowired RetryTemplate retryTemplate,
@Autowired ExecutorChannel kafkaErrorChannel) {
RequestHandlerRetryAdvice retryAdvice = new RequestHandlerRetryAdvice();
retryAdvice.setRecoveryCallback(new ErrorMessageSendingRecoverer(kafkaErrorChannel));
retryAdvice.setRetryTemplate(retryTemplate);
return retryAdvice;
}
@Bean
public RetryTemplate retryTemplate() {
return new RetryTemplateBuilder().maxAttempts(3).exponentialBackoff(1000, 3.0, 30000)
.retryOn(MessageHandlingException.class).build();
}
最后我有一个Mongo 处理程序,可以将失败的消息保存到某个集合中
@Bean
@ServiceActivator(inputChannel = "kafkaErrorChannel")
public MongoDbStoringMessageHandler kafkaFailureHandler(@Autowired MongoDatabaseFactory mongoDbFactory,
@Autowired MongoConverter mongoConverter) {
String collectionExpressionString = "headers['" + upstreamTypeHeader + "'] + '_'+ headers['" + upstreamInstanceHeader + "']+ '_FAIL'";
return getMongoDbStoringMessageHandler(mongoDbFactory, mongoConverter, collectionExpressionString);
}
我很难弄清楚我是否正确连接了所有这些,因为测试似乎从来没有工作过,在测试类中我不设置任何嵌入式 kafka 或连接到 kafka 所以该消息发布将失败,期望这会触发重试建议并最终保存到 mongo 中的死信集合。
@Test
void testFailedKafkaPublish() {
//Dummy message
Map<String, String> map = new HashMap<>();
map.put("key", "value");
// Publish Message
Message<Map<String, String>> message = MessageBuilder.withPayload(map)
.setHeader("X-UPSTREAM-TYPE", "alm")
.setHeader("X-INSTANCE-HEADER", "jira")
.build();
kafkaGateway.publish(message);
//assert successful message is saved in FAIL collection
assertThat(mongoTemplate.findAll(DBObject.class, "alm_jira_FAIL"))
.extracting("key")
.containsOnly("value");
}
我读到我们必须 setSync 到 Kafka 处理程序,以便它等待 kafka 操作的结果,所以我介绍了
@Value("${digite.swiftalk.kafka.test-mode:false}")
private boolean testMode;
到 Kafka 配置;在上面的测试中,我通过 @TestPropertySource 注释将其设置为 true:
@TestPropertySource(properties = {
"spring.main.banner-mode=off",
"spring.data.mongodb.database=swiftalk_db",
"spring.data.mongodb.port=29019",
"spring.data.mongodb.host=localhost",
"digite.swiftalk.kafka.test-mode=true",
})
我仍然看不到 Retry Advice 执行的任何注销或 Mongo 中保存的失败消息。另一个想法是使用Awaitility,但我不确定我应该在until() 方法中添加什么条件才能使其工作。
更新
为 Kafka 添加了调试日志,我注意到生产者进入了一个循环,试图在单独的线程中与 Kafka 重新连接
2021-03-25 10:56:02.640 DEBUG 66997 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=producer-1] Initiating connection to node localhost:9999 (id: -1 rack: null) using address localhost/127.0.0.1
2021-03-25 10:56:02.641 DEBUG 66997 --- [dPoolExecutor-1] o.a.k.clients.producer.KafkaProducer : [Producer clientId=producer-1] Kafka producer started
2021-03-25 10:56:02.666 DEBUG 66997 --- [ad | producer-1] o.apache.kafka.common.network.Selector : [Producer clientId=producer-1] Connection with localhost/127.0.0.1 disconnected
java.net.ConnectException: Connection refused
at java.base/sun.nio.ch.Net.pollConnect(Native Method) ~[na:na]
at java.base/sun.nio.ch.Net.pollConnectNow(Net.java:660) ~[na:na]
at java.base/sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:875) ~[na:na]
at org.apache.kafka.common.network.PlaintextTransportLayer.finishConnect(PlaintextTransportLayer.java:50) ~[kafka-clients-2.6.0.jar:na]
at org.apache.kafka.common.network.KafkaChannel.finishConnect(KafkaChannel.java:219) ~[kafka-clients-2.6.0.jar:na]
at org.apache.kafka.common.network.Selector.pollSelectionKeys(Selector.java:530) ~[kafka-clients-2.6.0.jar:na]
at org.apache.kafka.common.network.Selector.poll(Selector.java:485) ~[kafka-clients-2.6.0.jar:na]
at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:544) ~[kafka-clients-2.6.0.jar:na]
at org.apache.kafka.clients.producer.internals.Sender.runOnce(Sender.java:325) ~[kafka-clients-2.6.0.jar:na]
at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:240) ~[kafka-clients-2.6.0.jar:na]
at java.base/java.lang.Thread.run(Thread.java:832) ~[na:na]
当测试到达断言并因此失败
//assert successful message is saved in FAIL collection
assertThat(mongoTemplate.findAll(DBObject.class, "alm_jira_FAIL"))
.extracting("key")
.containsOnly("value");
所以看起来重试建议并没有首先接管前两次失败。
更新 2
更新配置类添加属性
@Value("${spring.kafka.producer.properties.max.block.ms:1000}")
private Integer productMaxBlockDurationMs;
并在kafkaTemplate配置方法中添加以下行
props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, productMaxBlockDurationMs);
修复了它。
更新 3
正如 Gary 所说,我们可以完全跳过添加所有这些道具等;我从课堂上删除了以下方法
@Bean
KafkaTemplate<String, String> kafkaTemplate() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, productMaxBlockDurationMs);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
}
并像这样从属性中简单地注入 kafka 配置,从而不必编写 Kafka 模板 bean
spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.producer.properties.max.block.ms=1000
spring.kafka.producer.properties.enable.idempotence=true
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
【问题讨论】:
标签: java spring-boot spring-integration spring-kafka