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

使用CompletableFuture发送Kafka消息时线程堆积及超时问题解决

问题解决方案

核心思路

  1. 利用Spring Kafka原生异步发送能力替代自定义@Async,避免线程池冲突堆积
  2. 每个消息发送后独立处理回调,不等待其他消息结果,彻底解决批次超时问题
  3. 合理配置Kafka生产者参数,优化批次行为
  4. 将@Retryable精准应用在单条消息的发送逻辑上,确保重试不影响全局流程

正确代码实现

1. Kafka生产者配置(application.yml)

spring:
  kafka:
    producer:
      bootstrap-servers: localhost:9092
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      # 调整批次参数,避免无意义的批次等待超时
      batch-size: 16384
      linger-ms: 5  # 最多等待5ms凑批次,无消息则立即发送
      request-timeout-ms: 30000
      retries: 3
      # 控制在途请求数,避免线程堆积
      properties:
        max.in.flight.requests.per.connection: 5

2. 审计日志发送服务

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.retry.annotation.Backoff;
import org.springframework.retry.annotation.Retryable;
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;

@Service
public class AuditLogKafkaService {

    private final KafkaTemplate<String, AuditLog> kafkaTemplate;
    private static final String AUDIT_TOPIC = "audit-logs";

    public AuditLogKafkaService(KafkaTemplate<String, AuditLog> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    // 仅对单条消息的发送失败进行重试,不阻塞其他消息
    @Retryable(
            retryFor = {Exception.class},
            backoff = @Backoff(delay = 1000, multiplier = 2), // 指数退避重试
            maxAttempts = 3
    )
    public CompletableFuture<Void> sendAuditLog(AuditLog auditLog) {
        // 使用KafkaTemplate原生异步发送,转换为CompletableFuture
        return kafkaTemplate.send(AUDIT_TOPIC, auditLog.getId(), auditLog)
                .completable()
                .whenComplete((result, ex) -> {
                    if (ex != null) {
                        // 记录失败日志,不要抛出异常(避免中断其他流程)
                        System.err.printf("审计日志[%s]发送失败: %s%n", auditLog.getId(), ex.getMessage());
                    } else {
                        SendResult<String, AuditLog> sendResult = (SendResult<String, AuditLog>) result;
                        System.out.printf("审计日志[%s]发送成功,offset: %d%n", auditLog.getId(), sendResult.getRecordMetadata().offset());
                    }
                })
                // 主线程无需等待发送完成,直接返回空结果
                .thenApply(ignored -> null);
    }
}

// 审计日志实体类
class AuditLog {
    private String id;
    private String transactionId;
    private String operation;

    // 构造器、Getter、Setter省略
    public String getId() { return id; }
    public void setId(String id) { this.id = id; }
    public String getTransactionId() { return transactionId; }
    public void setTransactionId(String transactionId) { this.transactionId = transactionId; }
    public String getOperation() { return operation; }
    public void setOperation(String operation) { this.operation = operation; }
}

3. 业务服务调用示例

import org.springframework.stereotype.Service;

@Service
public class TransactionService {

    private final AuditLogKafkaService auditLogService;

    public TransactionService(AuditLogKafkaService auditLogService) {
        this.auditLogService = auditLogService;
    }

    public void processTransaction(String transactionId) {
        // 执行业务核心逻辑
        System.out.println("处理事务: " + transactionId);

        // 构造审计日志
        AuditLog auditLog = new AuditLog();
        auditLog.setId(java.util.UUID.randomUUID().toString());
        auditLog.setTransactionId(transactionId);
        auditLog.setOperation("TRANSACTION_PROCESSED");

        // 发送审计日志,无需等待结果,不阻塞业务流程
        auditLogService.sendAuditLog(auditLog);

        // 继续执行后续业务逻辑
    }
}

关键说明

  • 移除@Async:KafkaTemplate的send()方法本身就是异步实现,无需额外添加@Async,避免自定义线程池与Spring Kafka内置线程池冲突导致线程堆积
  • 独立回调处理:每个消息的whenComplete仅处理自身的发送结果,不会等待其他消息,彻底解决批量等待引发的超时问题
  • Kafka参数优化:通过linger-ms控制批次等待时间,避免空等;max.in.flight.requests.per.connection限制在途请求数,防止线程过载
  • 精准重试:@Retryable仅作用于单条消息的发送逻辑,重试失败不会影响其他消息的发送流程
  • 非阻塞业务流程:业务代码无需调用join()或get()等待发送结果,确保事务处理效率不受Kafka发送影响

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 07:42:14