前言
在日常的微服务开发中,我们经常需要将核心接口的调用数据(如入参、返参、耗时等)异步上报到 Kafka,供下游进行数据同步、审计或数据计算。
如果在每个业务方法里手动编写 Kafka 发送逻辑,不仅代码冗余度高,而且严重侵入业务逻辑。本文将分享一种基于自定义注解 + AOP 拦截的优雅方案,实现接口调用数据的无侵入式、异步、防重上送,并支持线程池隔离与动态配置。
一、核心诉求与设计目标
在设计这套消息上送组件时,我们主要考虑了以下几个核心诉求:
| 诉求 | 说明 |
|---|---|
| 低侵入 | 业务代码尽量不感知 Kafka 上送逻辑,做到声明式调用 |
| 异步化 | 上送过程不能阻塞主业务线程,保障接口响应速度 |
| 防重复 | 同一请求(如同一分页数据)在同一天内避免重复上送 |
| 可配置 | Topic、线程池参数等可通过配置中心动态调整 |
| 链路追踪 | 异步线程需透传 TraceID,保证日志链路完整 |
二、整体架构设计
整体数据流转分为四层:业务层、AOP 拦截层、异步执行层和 Kafka 发送层。
┌─────────────────────────────────────────────────────────────┐ │ 业务 Service 层 │ │ @PushToKafka(scene="", type="") │ │ ↓ 方法调用 │ ├─────────────────────────────────────────────────────────────┤ │ AOP 拦截层 │ │ PushToKafkaAspect.around() │ │ ├─ 记录方法入参 / 返参 / 耗时 / 异常 │ │ ├─ Redis 防重锁校验 │ │ └─ 提交异步任务到线程池 │ ├─────────────────────────────────────────────────────────────┤ │ KafkaPushExecutorProvider │ │ ↓ 获取独立的 kafkaPushThreadPool │ ├─────────────────────────────────────────────────────────────┤ │ KafkaMessagePushProducerService │ │ ├─ 构造 KafkaMessageDto │ │ ├─ 注入公共字段 (如业务标识、接口名) │ │ └─ KafkaTemplate.send(ProducerRecord) │ ├─────────────────────────────────────────────────────────────┤ │ Kafka Broker │ └─────────────────────────────────────────────────────────────┘三、核心组件代码详解
3.1 自定义注解@PushToKafka
通过自定义注解标记需要上送 Kafka 的方法,实现业务代码零侵入。
packagecom.example.kafka.annotation;importjava.lang.annotation.*;/** * Kafka 推送注解 * <p> * 用于标记需要推送请求入参和接口返参到 Kafka 的方法。 * 通过 AOP 拦截该注解,异步将方法调用信息发送到 Kafka。 */@Target(ElementType.METHOD)@Retention(RetentionPolicy.RUNTIME)public@interfacePushToKafka{/** * 业务场景描述 */Stringscene()default"";/** * 业务类型标识 */Stringtype()default"";}设计要点:
@Target(ElementType.METHOD):仅作用于方法级别。@Retention(RetentionPolicy.RUNTIME):运行时保留,供 AOP 反射读取。type字段用于区分不同业务场景,最终写入 Kafka 消息体供下游分类消费。
3.2 AOP 拦截器PushToKafkaAspect
环绕通知拦截注解标记的方法,在方法执行完成后异步触发 Kafka 上送。
packagecom.example.kafka.aspect;importcom.example.common.model.dto.JsonResultDto;importcom.example.kafka.annotation.PushToKafka;importcom.example.kafka.config.KafkaPushExecutorProvider;importcom.example.kafka.service.KafkaMessagePushProducerService;importlombok.extern.slf4j.Slf4j;importorg.aspectj.lang.ProceedingJoinPoint;importorg.aspectj.lang.annotation.Around;importorg.aspectj.lang.annotation.Aspect;importorg.aspectj.lang.reflect.MethodSignature;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.beans.factory.annotation.Value;importorg.springframework.stereotype.Component;importjava.lang.reflect.Method;@Slf4j@Aspect@ComponentpublicclassPushToKafkaAspect{@AutowiredprivateKafkaMessagePushProducerServicekafkaProducerService;@AutowiredprivateKafkaPushExecutorProviderkafkaPushExecutorProvider;@Value("${app.kafka.topic.push-data:'datasync-public-logs'}")privateStringkafkaTopic;/** * 环绕通知:拦截 @PushToKafka 注解标记的方法 */@Around("@annotation(com.example.kafka.annotation.PushToKafka)")publicObjectaround(ProceedingJoinPointjoinPoint)throwsThrowable{longstartTime=System.currentTimeMillis();Objectresult=null;Throwableerror=null;try{result=joinPoint.proceed();returnresult;}catch(Throwablet){error=t;throwt;}finally{longcostTime=System.currentTimeMillis()-startTime;asyncPushToKafka(joinPoint,result,error,costTime);}}/** * 异步推送数据到 Kafka */privatevoidasyncPushToKafka(ProceedingJoinPointjoinPoint,Objectresult,Throwableerror,longcostTime){try{MethodSignaturesignature=(MethodSignature)joinPoint.getSignature();Methodmethod=signature.getMethod();PushToKafkapushToKafka=method.getAnnotation(PushToKafka.class);if(pushToKafka==null)return;Object[]args=joinPoint.getArgs();// 此处可根据实际业务校验入参和返参类型,例如:// if (args == null || args.length == 0 || !(args[0] instanceof BaseReq)) return;// if (result == null || !(result instanceof JsonResultDto)) return;// 获取业务参数StringinterfaceName=joinPoint.getSignature().getName();Stringtype=pushToKafka.type();log.info("开始异步推送 Kafka 数据,type: {}, interfaceName: {}, costTime: {}ms",type,interfaceName,costTime);// 提交异步任务到独立线程池kafkaPushExecutorProvider.getExecutor().execute(()->kafkaProducerService.sendMessage(kafkaTopic,type,interfaceName,result));}catch(Exceptione){log.error("异步推送 Kafka 数据异常",e);}}}核心逻辑拆解:
- 环绕拦截:
@Around拦截目标方法,记录startTime。 - 执行原方法:
joinPoint.proceed(),异常原样抛出,不影响主流程。 - 异步提交:通过专用线程池执行 Kafka 发送,避免阻塞主线程。
3.3 Kafka 生产者服务与 Redis 防重锁
负责消息序列化、字段注入、防重校验及最终发送到 Kafka Broker。
packagecom.example.kafka.service;importcom.alibaba.fastjson.JSON;importlombok.extern.slf4j.Slf4j;importorg.apache.kafka.clients.producer.ProducerRecord;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.data.redis.core.RedisTemplate;importorg.springframework.kafka.core.KafkaTemplate;importorg.springframework.kafka.support.SendResult;importorg.springframework.stereotype.Service;importorg.springframework.util.concurrent.ListenableFuture;importorg.springframework.util.concurrent.ListenableFutureCallback;importjava.util.Calendar;importjava.util.concurrent.TimeUnit;@Slf4j@ServicepublicclassKafkaMessagePushProducerService{@AutowiredprivateKafkaTemplate<String,String>kafkaTemplate;@AutowiredprivateRedisTemplate<String,String>redisTemplate;/** * 发送消息到 Kafka */publicvoidsendMessage(Stringtopic,Stringtype,StringinterfaceName,Objectresult){// 1. 构造消息体并序列化StringmessageJson=JSON.toJSONString(result);// 2. 构建 ProducerRecordProducerRecord<String,String>record=newProducerRecord<>(topic,messageJson);// 3. 异步发送并监听回调ListenableFuture<SendResult<String,String>>future=kafkaTemplate.send(record);future.addCallback(newListenableFutureCallback<SendResult<String,String>>(){@OverridepublicvoidonSuccess(SendResult<String,String>result){log.debug("Kafka 消息发送成功, topic={}",topic);}@OverridepublicvoidonFailure(Throwableex){log.error("Kafka 消息发送失败, topic={}",topic,ex);}});}/** * 检查 Redis 锁是否存在 (防重机制) */publicbooleanisLock(StringlockKey){returnBoolean.TRUE.equals(redisTemplate.hasKey(lockKey));}/** * 设置 Redis 锁,过期时间为当天最后一分钟 */publicvoidlock(StringlockKey){Calendarnow=Calendar.getInstance();CalendarendOfDay=Calendar.getInstance();endOfDay.set(Calendar.HOUR_OF_DAY,23);endOfDay.set(Calendar.MINUTE,59);endOfDay.set(Calendar.SECOND,59);longexpireMinutes=TimeUnit.MILLISECONDS.toMinutes(endOfDay.getTimeInMillis()-now.getTimeInMillis())+1;if(expireMinutes<=0)expireMinutes=1;redisTemplate.opsForValue().set(lockKey,"1",expireMinutes,TimeUnit.MINUTES);}}防重锁特性:
- 粒度细:以
接口名 + 业务主键 + 分页参数为维度,避免同一分页数据重复上送。 - 当天有效:过期时间动态计算到当天
23:59:59,次日自动失效并可重新上送。
3.4 独立线程池配置
Kafka 上送采用独立线程池,与业务线程隔离。推荐使用ThreadPoolTaskExecutor并结合链路追踪组件(如 Sleuth/Micrometer)透传 TraceID。
# application.yml (支持配置中心动态下发)app:trace:thread:configList:-threadNameBean:kafkaPushThreadPoolcorePoolSize:4maxPoolSize:16queueCapacity:200keepAliveSeconds:60threadNamePrefix:kafka-push-# 队列满时由调用线程自己执行,避免丢数据rejectedExecutionHandler:java.util.concurrent.ThreadPoolExecutor$CallerRunsPolicy四、使用方式
在需要上送 Kafka 的 Service 方法上添加@PushToKafka注解即可,真正做到了声明式编程:
@ServicepublicclassDataService{@PushToKafka(scene="核心业务数据同步",type="business_data_sync")publicJsonResultDto<DataVO>queryData(QueryReqreq){// 1. 纯业务逻辑...DataVOdata=doBusinessLogic(req);returnJsonResultDto.success(data);}}零侵入体现:
- ❌ 无需修改方法内部逻辑
- ❌ 无需手动构造 Kafka 消息
- ❌ 无需关注线程池和异常处理
五、关键设计要点总结
| 设计点 | 方案 | 收益 |
|---|---|---|
| 注解驱动 | 自定义@PushToKafka+ AOP 拦截 | 业务代码零侵入,声明式使用 |
| 异步发送 | 独立线程池kafkaPushThreadPool | 不阻塞主业务线程,保障接口 RT |
| 防重机制 | Redis 分布式锁,当天维度去重 | 避免重复数据,降低 Kafka 压力 |
| 失败感知 | ListenableFutureCallback回调 | 发送失败可记录日志,便于排查 |
| 动态配置 | Topic、线程池均走 Nacos/Apollo | 无需发版即可调整参数 |
| 链路追踪 | 支持 TraceID 透传的线程池包装 | TraceID 在线程池间透传,日志可串联 |
六、生产环境注意事项
- 切面执行顺序:若方法上存在其他切面(如
@Transactional、@Cacheable),需关注@Order优先级,确保 Kafka 切面在最外层,能捕获到最终的完整结果。 - 序列化兼容性:消息体建议统一转为 JSON 字符串,并与下游消费方约定好字段结构,避免反序列化失败。
- Redis 锁失效兜底:若 Redis 出现网络抖动导致锁未正常写入,可能会产生少量重复消息。必须在 Kafka 消费端做好幂等性处理(如基于唯一键去重)。
- 线程池监控:建议对
kafkaPushThreadPool的活跃线程数、队列积压量配置 Prometheus 监控与告警,防止队列打满触发拒绝策略。
原创不易,如果这篇文章对你有帮助,欢迎点赞、收藏、关注!你的支持是我持续创作的动力。