【问题标题】:How to get RetryAdvice working for KafkaProducerMessageHandler如何让 RetryAdvice 为 KafkaProducerMessageHandler 工作
【发布时间】: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


    【解决方案1】:

    KafkaProducers 在失败前默认阻塞 60 秒。

    尝试减少max.block.ms producer 属性。

    https://kafka.apache.org/documentation/#producerconfigs_max.block.ms

    编辑

    这是一个例子:

    @SpringBootApplication
    public class So66768745Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So66768745Application.class, args);
        }
    
        @Bean
        IntegrationFlow flow(KafkaTemplate<String, String> template, RequestHandlerRetryAdvice retryAdvice) {
            return IntegrationFlows.from(Gate.class)
                    .handle(Kafka.outboundChannelAdapter(template)
                                .topic("testTopic"), e -> e
                            .advice(retryAdvice))
                    .get();
        }
    
        @Bean
        RequestHandlerRetryAdvice retryAdvice(QueueChannel channel) {
            RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice();
            advice.setRecoveryCallback(new ErrorMessageSendingRecoverer(channel));
            return advice;
        }
    
        @Bean
        QueueChannel channel() {
            return new QueueChannel();
        }
    
    }
    
    interface Gate {
    
        void sendToKafka(String out);
    
    }
    
    @SpringBootTest
    @TestPropertySource(properties = {
            "spring.kafka.bootstrap-servers: localhost:9999",
            "spring.kafka.producer.properties.max.block.ms: 500" })
    class So66768745ApplicationTests {
    
        @Autowired
        Gate gate;
    
        @Autowired
        QueueChannel channel;
    
        @Test
        void test() {
            this.gate.sendToKafka("test");
            Message<?> em = this.channel.receive(60_000);
            assertThat(em).isNotNull();
            System.out.println(em);
        }
    
    }
    
    2021-03-23 15:16:13.908 ERROR 2668 --- [           main] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='null' and payload='test' to topic testTopic:
    
    org.apache.kafka.common.errors.TimeoutException: Topic testTopic not present in metadata after 500 ms.
    
    2021-03-23 15:16:14.343  WARN 2668 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient   : [Producer clientId=producer-1] Connection to node -1 (localhost/127.0.0.1:9999) could not be established. Broker may not be available.
    2021-03-23 15:16:14.343  WARN 2668 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient   : [Producer clientId=producer-1] Bootstrap broker localhost:9999 (id: -1 rack: null) disconnected
    2021-03-23 15:16:14.415 ERROR 2668 --- [           main] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='null' and payload='test' to topic testTopic:
    
    org.apache.kafka.common.errors.TimeoutException: Topic testTopic not present in metadata after 500 ms.
    
    2021-03-23 15:16:14.921 ERROR 2668 --- [           main] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='null' and payload='test' to topic testTopic:
    
    org.apache.kafka.common.errors.TimeoutException: Topic testTopic not present in metadata after 500 ms.
    
    ErrorMessage [payload=org.springframework.messaging.MessagingException: Failed to handle; nested exception is org.springframework.kafka.KafkaException: Send failed; nested exception is org.apache.kafka.common.errors.TimeoutException: Topic testTopic not present in metadata after 500 ms., failedMessage=GenericMessage [payload=test, headers={replyChannel=nullChannel, errorChannel=, id=d8ce277a-3d9a-b0bc-c14b-80d63ca13858, timestamp=1616526973218}], headers={id=1a6c29d2-f8d8-adf0-7569-db7610b020ef, timestamp=1616526974921}]
    

    【讨论】:

    • 我添加了一个有效的示例。请注意,对于这个特定错误(无法获取元数据),不需要同步,因为在调用线程上抛出了 TimeoutException。对于其他(异步)错误,您将需要 sync
    • 谢谢! @Gary Russel,尝试了您的更改,也意识到我可以编写流程以使其清晰:-) 但我仍然无法复制您的答案中的行为;你能帮我指出我的代码哪里出错了吗?这是要点gist.github.com/anadimisra/087983f7eb3fe1d814c91adf9c54f4ce
    • 很难说;尝试将 org.apache.kafka 日志级别设置为 DEBUG 以查看它是否提供任何线索 - 确保在生产配置信息日志中应用 max.block.ms。如果您仍然无法弄清楚,请将其剥离到最低限度(如我的)并将完整的项目发布到某个地方,以便我可以在本地运行它。我将示例更改为使用执行程序通道,它仍然有效。
    • 对,但这是生产者 I/O 线程,main 线程应该像我的示例中那样超时;确保max.block.ms 配置正确(在ProducerConfig INFO 日志中)。如果您仍然无法弄清楚,请在某个地方发布一个最小的项目(如我的),我会看看有什么问题。
    • application.yml/properties 仅自动应用于 Boot 的自动配置 bean(在这种情况下为生产者工厂)。如果您定义自己的基础设施 bean,则必须完全配置它们。一般不需要自己定义生产者工厂bean,为什么不直接使用Boot的自动配置工厂呢?
    猜你喜欢
    • 2021-11-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-04-11
    • 2020-12-12
    • 1970-01-01
    相关资源
    最近更新 更多