You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将Kafka阻塞与非阻塞重试记录为警告日志?

Kafka阻塞式与非阻塞式重试的警告日志记录方案

一、阻塞式重试日志记录

阻塞式重试一般基于同步发送(如producer.send(record).get()),可以通过手动封装重试逻辑+异常捕获实现精准的日志控制:

  • 核心思路:在重试循环内,捕获可重试异常(RetriableException及其子类),每次重试时记录警告日志,包含异常栈、当前重试次数、目标主题等关键信息。
  • 代码示例:
private static final int MAX_RETRIES = 3;
private final Logger logger = LoggerFactory.getLogger(YourProducerClass.class);

public void sendWithBlockingRetry(ProducerRecord<String, String> record) throws Exception {
    int retryCount = 0;
    while (true) {
        try {
            producer.send(record).get();
            logger.info("消息发送成功,主题:{},分区:{}", record.topic(), record.partition());
            break;
        } catch (ExecutionException e) {
            Throwable cause = e.getCause();
            // 仅针对可重试异常触发重试与日志记录
            if (cause instanceof RetriableException && retryCount < MAX_RETRIES) {
                retryCount++;
                logger.warn("[阻塞式重试] 消息发送失败,将进行第{}次重试。异常详情:{}", 
                    retryCount, 
                    ExceptionUtils.getStackTrace(cause));
                // 可选:添加指数退避间隔
                Thread.sleep(1000 * retryCount);
            } else {
                logger.error("消息发送最终失败,超出最大重试次数或为不可重试异常。主题:{},异常详情:{}",
                    record.topic(),
                    ExceptionUtils.getStackTrace(cause));
                throw e;
            }
        }
    }
}

二、非阻塞式重试日志记录

非阻塞式重试依赖异步回调或重试框架,核心是在异常触发的直接拦截点记录日志,无需依赖ProducerInterceptor:

方案1:自定义异步发送Callback

直接在KafkaProducer.send()的回调中处理异常与重试逻辑,同步记录日志:

private void sendWithNonBlockingRetry(ProducerRecord<String, String> record) {
    int[] retryCount = {0}; // 用数组传递可变值,适配Lambda的不可变变量限制
    sendAsync(record, retryCount);
}

private void sendAsync(ProducerRecord<String, String> record, int[] retryCount) {
    producer.send(record, (metadata, exception) -> {
        if (exception != null) {
            Throwable cause = exception.getCause() != null ? exception.getCause() : exception;
            if (cause instanceof RetriableException && retryCount[0] < MAX_RETRIES) {
                retryCount[0]++;
                logger.warn("[非阻塞式重试] 消息发送失败,将进行第{}次重试。主题:{},异常详情:{}",
                    retryCount[0],
                    record.topic(),
                    ExceptionUtils.getStackTrace(cause));
                // 递归触发下一次异步重试
                sendAsync(record, retryCount);
            } else {
                logger.error("消息发送最终失败,主题:{},异常详情:{}",
                    record.topic(),
                    ExceptionUtils.getStackTrace(cause));
            }
        } else {
            logger.info("消息异步发送成功,主题:{},分区:{},偏移量:{}",
                metadata.topic(),
                metadata.partition(),
                metadata.offset());
        }
    });
}

方案2:结合Spring Retry与KafkaTemplate(Spring生态)

如果使用Spring Kafka,可通过RetryTemplate配合自定义RetryListener统一拦截重试事件,记录警告日志:

@Configuration
public class KafkaRetryConfig {
    private final Logger logger = LoggerFactory.getLogger(KafkaRetryConfig.class);
    private static final int MAX_RETRIES = 3;

    @Bean
    public RetryTemplate kafkaRetryTemplate() {
        RetryTemplate retryTemplate = new RetryTemplate();
        // 配置重试策略:仅重试可重试异常,最多3次
        SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
        retryPolicy.setMaxAttempts(MAX_RETRIES);
        Map<Class<? extends Throwable>, Boolean> retryableExceptions = new HashMap<>();
        retryableExceptions.put(RetriableException.class, true);
        retryPolicy.setRetryableExceptions(retryableExceptions);
        retryTemplate.setRetryPolicy(retryPolicy);

        // 注册重试监听器,捕获重试事件并记录日志
        retryTemplate.registerListener(new RetryListener() {
            @Override
            public <T, E extends Throwable> boolean open(RetryContext context, RetryCallback<T, E> callback) {
                return true;
            }

            @Override
            public <T, E extends Throwable> void close(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {}

            @Override
            public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {
                int retryCount = context.getRetryCount();
                String topic = (String) context.getAttribute("topic");
                logger.warn("[Spring Retry 非阻塞重试] 第{}次重试触发,主题:{},异常详情:{}",
                    retryCount,
                    topic,
                    ExceptionUtils.getStackTrace(throwable));
            }
        });
        return retryTemplate;
    }
}

// 使用示例
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
@Autowired
private RetryTemplate kafkaRetryTemplate;

public void sendWithSpringRetry(ProducerRecord<String, String> record) {
    kafkaRetryTemplate.execute(context -> {
        context.setAttribute("topic", record.topic());
        kafkaTemplate.send(record).get(); // 同步调用让RetryTemplate捕获异常,可结合异步回调调整
        return null;
    });
}

关键注意事项

  • 仅对可重试异常(RetriableException)记录重试日志,避免对无效异常(如InvalidTopicException)生成冗余日志。
  • 日志需包含重试次数、目标主题、完整异常栈,便于后续问题排查。
  • 非阻塞场景下注意线程安全,如示例中用数组传递重试计数,避免Lambda的变量不可变限制。

内容的提问来源于stack exchange,提问作者Casper T

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 09:53:12