kafkaTemplate;
/**
* 同步发送消息
* @param topic 主体
* @param msg 消息
*/
public void sendMessageSync(String topic, Object msg) throws ExecutionException, InterruptedException {
try {
kafkaTemplate.send(topic, msg).get();
log.info("sendMessageSync 同步消息发送 成功 topic={} msg={}", topic, msg);
} catch (Exception e) {
log.error("sendMessageSync 同步消息发送 失败", e);
throw e;
}
}
/**
* 异步发送消息
* @param topic 主体
* @param msg 消息
*/
public void sendMessageAsync(String topic, Object msg) {
kafkaTemplate.send(topic, msg).addCallback(success -> {
log.info("sendMessageAsync 异步消息发送 成功 topic={} msg={}", topic, msg);
}, failure -> {
log.error("sendMessageAsync 异步消息发送 失败", failure);
});
}
}
```
## kafka 消费者:KafkaConsumerService.java
```java
/**
* kafka 消费者
*/
@Service
@Slf4j
public class KafkaConsumerService {
/**
* 监听Kafka消息 手动提交
* todo 测试topic 需修改
*
* @param msg 消息内容
*/
@KafkaListener(topics = "test-reject", groupId = "dynamic-pool-service")
public void consume(String msg, Acknowledgment ack) {
log.info("kafka 开始消费 msg: {}", msg);
try {
// 这里是通过反射 进行任务执行
KafkaRejectPolicy.RejectedTaskMessage.taskStart(msg);
ack.acknowledge();
} catch (Exception e) {
log.info("kafka 消费异常", e);
}
}
}
```
## kafka 拒绝策略
> 这里面重要的就是 KafkaTaskInfoRunnable 这个抽象类
>
> 我们的任务都需要继承这个类来实现 kafka 异步任务来做到任务不丢失
>
> 逻辑是这样的(后面测试案例会讲到,这里只是给大家理一下逻辑)
>
> 我们先实现一个 A 类继承 KafkaTaskInfoRunnable 这个抽象类(A.java extends KafkaTaskInfoRunnable.java)
>
> KafkaTaskInfoRunnable 实现了 Runnable 接口,所以 A 类需要实现 run 方法,这样就可以通过线程池进行异步任务
>
> 我们线程池满了之后,任务被拒绝了,会给 kafka 发送一个消息(json字符串)
>
> json 里面包括了实现任务的类(A.java)(Class)、入参(具体值)、入参类型(Class)等(这三个比较关键)
>
> kafka 监听后会从消息(json)里面获取到任务的类(Class)、入参(具体值)、入参类型(Class)
>
> 然后通过反射同步的调用 A.java 的 run 方法,达到消峰(因为kafka会一个任务执行完再接着执行)和任务不丢失
>
> 由于要使用反射获取 A.java 的具体数据,所以 A.java 是不能使用匿名内部类的
>
> ps:发送警告我没有处理,有兴趣大家可以实现以下,这个也简单,如果警告频繁,说明 QPS 过大,我们就需要自己调整线程池或者就是加机器了
```java
/**
* Kafka拒绝策略 - 将被拒绝的任务发送到Kafka队列
*/
@Slf4j
public class KafkaRejectPolicy implements RejectedExecutionHandler {
private static final AtomicLong rejectCount = new AtomicLong(0);
private final String topic;
/**
* 应用名称(用于消息标识)
*/
private final String applicationName;
/**
* 构造函数
*
* @param topic 目标主题
*/
public KafkaRejectPolicy(String topic) {
this.topic = topic;
// 这里得先在启动器配置优先加载 SpringBootBeanUtil 的 bean 对象
String applicationName = SpringBootBeanUtil.getEnvironment().getProperty("spring.application.name");
if (applicationName == null || applicationName.isEmpty()) {
throw new IllegalArgumentException("配置项 spring.application.name 不能为空");
}
this.applicationName = applicationName;
}
/**
* 获取ip 地址
*/
private static String getIpAddress() throws UnknownHostException {
InetAddress localHost = InetAddress.getLocalHost();
return localHost.getHostAddress();
}
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
try {
// 检查任务类型
if (!(r instanceof KafkaTaskInfoRunnable)) {
// 使用 submit 方法时会走到 AbstractExecutorService 的 submit 方法,会封装一层 FutureTask
// 如果是 FutureTask,尝试提取内部的 Runnable
if (r instanceof FutureTask) {
try {
// 通过反射获取 FutureTask 内部的 callable
Field callableField = FutureTask.class.getDeclaredField("callable");
callableField.setAccessible(true);
Object callable = callableField.get(r);
// 如果是 RunnableAdapter(Executor.callable 返回的内部类)
if (callable != null && callable.getClass().getName().contains("RunnableAdapter")) {
Field taskField = callable.getClass().getDeclaredField("task");
taskField.setAccessible(true);
Object originalTask = taskField.get(callable);
if (originalTask instanceof KafkaTaskInfoRunnable) {
r = (Runnable) originalTask;
}
}
} catch (Exception e) {
log.error("提取 FutureTask 内部任务失败", e);
}
}
}
// 最终检查
if (!(r instanceof KafkaTaskInfoRunnable)) {
throw new IllegalArgumentException("该任务无法由 Kafka 执行,必须使用 KafkaTaskInfoRunnable 类型");
}
KafkaTaskInfoRunnable> KafkaTaskInfoRunnable = (KafkaTaskInfoRunnable>) r;
long count = rejectCount.incrementAndGet();
// 获取线程池状态信息
String poolInfo = getThreadPoolInfo(executor);
// 构建拒绝任务的消息
RejectedTaskMessage taskMessage = new RejectedTaskMessage()
.setTaskId(generateTaskId())
.setApplicationName(applicationName)
.setIpAddress(getIpAddress())
.setThreadPoolInfo(poolInfo)
.setRejectTime(System.currentTimeMillis())
.setRejectCount(count)
.setTaskDataJson(KafkaTaskInfoRunnable.getTaskDataJson())
.setTaskClass(KafkaTaskInfoRunnable.getTaskClass())
.setTaskInputParameterClass(KafkaTaskInfoRunnable.getTaskInputParameterClass());
// 发送到Kafka
String msg = JSONObject.toJSONString(taskMessage);
// 使用 kafka 生产者发送消息
KafkaProducerService kafkaProducer = SpringBootBeanUtil.getBean(KafkaProducerService.class);
// 同步发送并处理回调
kafkaProducer.sendMessageSync(topic, msg);
log.info("拒绝任务已发送到Kafka队列。拒绝次数: {}, 线程池状态: {}", count, poolInfo);
} catch (Exception e) {
log.error(" kafka 拒绝策略 发送消息异常", e);
// 降级处理(让当前线程自己处理)
if (!executor.isShutdown()) {
r.run();
}
}
}
/**
* 获取线程池状态信息
*
* 包含以下关键指标:
* - 活跃线程数:当前正在执行任务的线程数量
* - 线程池大小:线程池中当前的线程数量(包括空闲线程)
* - 核心线程数:线程池保持的最小线程数量
* - 最大线程数:线程池允许创建的最大线程数量
* - 队列大小:任务队列中等待执行的任务数量
* - 队列容量:任务队列的最大容量
*/
private String getThreadPoolInfo(ThreadPoolExecutor executor) {
BlockingQueue queue = executor.getQueue();
return String.format(
"活跃线程: %d, 池中线程: %d, 核心线程: %d, 最大线程: %d, 队列大小/队列容量: %d/%d",
// 正在执行任务的线程数
executor.getActiveCount(),
// 线程池中当前的线程总数
executor.getPoolSize(),
// 核心线程数
executor.getCorePoolSize(),
// 最大线程数
executor.getMaximumPoolSize(),
// 队列中等待的任务数
queue.size(),
// 队列总容量
queue.size() + queue.remainingCapacity()
);
}
/**
* 生成任务ID
*/
private String generateTaskId() {
return UUID.randomUUID().toString();
}
/**
* 发送告警
*/
private void sendAlert(RejectedTaskMessage message) {
// 实现告警逻辑(邮件、短信、钉钉等)
log.error("任务处理失败!任务ID: {}", message.getTaskId());
}
/**
* 任务信息接口(不能使用匿名内部类)
* 实现此接口的任务可以传递更多信息到拒绝策略
*/
public static abstract class KafkaTaskInfoRunnable implements Runnable {
private final T taskData;
public KafkaTaskInfoRunnable(T taskData) {
// 运行时检查
if (this.getClass().isAnonymousClass()) {
throw new IllegalStateException("TaskInfoRunnable 类 不能使用匿名内部类");
}
this.taskData = taskData;
}
protected String getTaskDataJson() {
return JSONObject.toJSONString(taskData);
}
protected T getTaskData() {
return this.taskData;
}
/**
* 任务执行类名(全类名)
*/
protected Class> getTaskClass() {
return (Class>) this.getClass();
}
/**
* 任务入参类名(全类名)
*/
protected Class getTaskInputParameterClass() {
return (Class) taskData.getClass();
}
}
/**
* 拒绝任务消息实体
*/
@Data
@Accessors(chain = true)
public static class RejectedTaskMessage implements Serializable {
private static final long serialVersionUID = 1L;
/**
* 任务ID(唯一ID)
*/
private String taskId;
/**
* 应用名称
*/
private String applicationName;
/**
* ip地址
*/
private String ipAddress;
/**
* 线程池状态信息
*/
private String threadPoolInfo;
/**
* 拒绝时间戳
*/
private long rejectTime;
/**
* 拒绝计数
*/
private long rejectCount;
/**
* 任务执行类名(全类名)
*/
private Class extends KafkaTaskInfoRunnable> taskClass;
/**
* 任务入参类名(全类名)
*/
private Class> taskInputParameterClass;
/**
* 任务入参数据(json)
*/
private String taskDataJson;
public static void taskStart(String msgJson) throws Exception {
JSON.config(JSONReader.Feature.SupportClassForName);
KafkaRejectPolicy.RejectedTaskMessage message = JSONObject.parseObject(
msgJson,
KafkaRejectPolicy.RejectedTaskMessage.class
);
String dataJson = message.getTaskDataJson();
Class> inputParameterClass = message.getTaskInputParameterClass();
Class> taskClass = (Class>) message.getTaskClass();
Object o = JSONObject.parseObject(dataJson, inputParameterClass);
Constructor> constructor = taskClass.getDeclaredConstructor(inputParameterClass);
constructor.setAccessible(true);
KafkaTaskInfoRunnable> infoRunnable = constructor.newInstance(o);
infoRunnable.run();
}
}
}
```
## 动态线程池:DynamicKafkaRejectThreadPoolConfig.java
> 跟上面的动态线程池类似,但有细微差距
>
> 最主要的是 ThreadPoolKafkaTaskExecutor 类
>
> 该类也继承了 ThreadPoolTaskExecutor 类,用于可以动态改变线程池的参数(核心、最大线程、队列容量(也是通过VariableLinkedBlockingQueue 实现的))
>
> 这里比较重要的就是重写线程池的 execute 和 submit 方法,只允许继承了 KafkaTaskInfoRunnable 的任务入队
```java
/**
* 动态线程池 (kafka 拒绝策略版本)
*
* 加 todo 的 就是可以选择性修改的
*/
@Configuration
@RefreshScope // 支持配置热更新
@EnableAsync // 启用异步支持
@Slf4j
@Lazy
public class DynamicKafkaRejectThreadPoolConfig {
/**
* 核心参数 默认使用io密集型
* todo 根据业务情况自行修改
*/
@Value("${dynamic.thread-pool.core-size:0}")
private int coreSize;
/**
* 最大参数 默认使用io密集型
* todo 根据业务情况自行修改
*/
@Value("${dynamic.thread-pool.max-size:0}")
private int maxSize;
/**
* 队列容量 默认 1w 的队列
* todo 根据业务情况自行修改
*/
@Value("${dynamic.thread-pool.queue-capacity:10000}")
private int queueCapacity;
/**
* 超时时间 默认 60s 超时
* todo 根据业务情况自行修改
*/
@Value("${dynamic.thread-pool.keep-alive-seconds:60}")
private int keepAliveSeconds;
/**
* 超时时间 默认 60s 超时
* todo 根据业务情况自行修改
*/
@Value("${dynamic.thread-pool.kafka-reject-policy-topic:test-reject}")
private String kafkaRejectPolicyTopic;
/**
* 动态线程池
* todo 根据业务情况自行修改 Bean 的值
*/
@Bean("dynamicKafkaRejectThreadPool")
ThreadPoolKafkaTaskExecutor dynamicKafkaRejectThreadPool() {
int cpuCores = Runtime.getRuntime().availableProcessors();
if (coreSize <= 0) {
coreSize = cpuCores * 2;
}
if (maxSize <= 0) {
maxSize = cpuCores * 4;
}
ThreadPoolKafkaTaskExecutor executor = new ThreadPoolKafkaTaskExecutor();
// 核心线程数(默认线程数)
executor.setCorePoolSize(coreSize);
// 最大线程数(队列满后扩容的最大值)
executor.setMaxPoolSize(maxSize);
// 线程空闲时间(秒)
executor.setKeepAliveSeconds(keepAliveSeconds);
// 队列容量(超过核心线程数时,任务进入队列)
executor.setQueueCapacity(queueCapacity);
// 线程名前缀
executor.setThreadNamePrefix("dynamic-thread-");
// 拒绝策略(此处使用报错执行)
executor.setRejectedExecutionHandler(new KafkaRejectPolicy(kafkaRejectPolicyTopic));
// 初始化
executor.initialize();
return executor;
}
/**
* 动态线程池
* todo 如果你改了 Bean 的值 记得改 Qualifier 中的值
*/
@Autowired
@Qualifier("dynamicKafkaRejectThreadPool")
@Lazy
private ThreadPoolTaskExecutor dynamicKafkaRejectThreadPool;
/**
* 公平读写锁
*/
private final static ReentrantReadWriteLock reentrantLock = new ReentrantReadWriteLock(true);
/**
* 动态调整线程池参数
*
* @param corePoolSize 核心线程数
* @param maxPoolSize 最大线程数
* @param keepAliveSeconds 超时时间(秒)
* @param queueCapacity 队列容量
*/
public void adjustThreadPool(Integer corePoolSize, Integer maxPoolSize, Integer keepAliveSeconds, Integer queueCapacity) {
ReentrantReadWriteLock.WriteLock writeLock = reentrantLock.writeLock();
try {
writeLock.lock();
// 调整核心线程数
if (corePoolSize != null && corePoolSize != dynamicKafkaRejectThreadPool.getCorePoolSize()) {
dynamicKafkaRejectThreadPool.setCorePoolSize(corePoolSize);
}
// 调整最大线程数
if (maxPoolSize != null && maxPoolSize != dynamicKafkaRejectThreadPool.getMaxPoolSize()) {
dynamicKafkaRejectThreadPool.setMaxPoolSize(maxPoolSize);
}
// 调整超时时间
if (keepAliveSeconds != null && keepAliveSeconds != dynamicKafkaRejectThreadPool.getKeepAliveSeconds()) {
dynamicKafkaRejectThreadPool.setKeepAliveSeconds(keepAliveSeconds);
}
// 调整队列容量(需重置队列)
if (queueCapacity != null && queueCapacity != dynamicKafkaRejectThreadPool.getQueueCapacity()) {
VariableLinkedBlockingQueue queue =
(VariableLinkedBlockingQueue) dynamicKafkaRejectThreadPool.getThreadPoolExecutor().getQueue();
if (queueCapacity < queue.remainingCapacity() + queue.size()) {
throw new UnsupportedOperationException("线程池 不支持缩容");
}
// 需要自定义可调整队列容量
queue.setCapacity(queueCapacity);
// 重置线程池容量参数(这里只是改参数,用于查询时与队列保持同步,并不能调整队列大小)
dynamicKafkaRejectThreadPool.setQueueCapacity(queueCapacity);
}
} finally {
writeLock.unlock();
}
}
/**
* 获取线程池状态
*/
public ThreadPoolStatus getThreadPoolStatus() {
ReentrantReadWriteLock.ReadLock readLock = reentrantLock.readLock();
try {
readLock.lock();
ThreadPoolExecutor threadPoolExecutor = dynamicKafkaRejectThreadPool.getThreadPoolExecutor();
BlockingQueue queue = threadPoolExecutor.getQueue();
int useQueueCapacity = queue.size();
int remainingCapacity = queue.remainingCapacity();
return new ThreadPoolStatus()
.setCorePoolSize(dynamicKafkaRejectThreadPool.getCorePoolSize())
.setMaxPoolSize(dynamicKafkaRejectThreadPool.getMaxPoolSize())
.setActiveThreads(dynamicKafkaRejectThreadPool.getActiveCount())
.setKeepAliveSeconds(dynamicKafkaRejectThreadPool.getKeepAliveSeconds())
.setQueueMaxCapacity(useQueueCapacity + remainingCapacity)
.setQueueUseCapacity(useQueueCapacity)
.setQueueRemainingCapacity(remainingCapacity);
} finally {
readLock.unlock();
}
}
/**
* 监听配置文件变更
* todo 监听配置文件变更
*/
@NacosConfigKeysListener(dataId = "dynamic-pool-service.yaml",
group = "DEFAULT_GROUP", interestedKeyPrefixes = "dynamic.thread-pool.")
private void onConfigChanged(ConfigChangeEvent changeEvent) {
log.info("onConfigChanged 监听到配置变化:{}", JSONObject.toJSONString(changeEvent));
DynamicPoolConfig config = new DynamicPoolConfig();
try {
Class extends DynamicPoolConfig> aClass = config.getClass();
Field[] declaredFields = aClass.getDeclaredFields();
for (Field declaredField : declaredFields) {
declaredField.setAccessible(true);
if (declaredField.isAnnotationPresent(ConfigChangeEventToBean.class)) {
ConfigChangeEventToBean annotation = declaredField.getAnnotation(ConfigChangeEventToBean.class);
String value = annotation.value();
ConfigChangeItem changeItem = changeEvent.getChangeItem(value);
// 处理配置变更
if (changeItem != null && PropertyChangeType.MODIFIED.equals(changeItem.getType())) {
declaredField.set(config, changeItem.getNewValue());
}
}
}
log.info("onConfigChanged 配置转换结果:{}", config);
this.adjustThreadPool(
Optional.ofNullable(config.getCoreSize()).map(Integer::parseInt).orElse(null),
Optional.ofNullable(config.getMaxSize()).map(Integer::parseInt).orElse(null),
Optional.ofNullable(config.getKeepAliveSeconds()).map(Integer::parseInt).orElse(null),
Optional.ofNullable(config.getQueueCapacity()).map(Integer::parseInt).orElse(null)
);
} catch (Exception e) {
log.error("onConfigChanged 配置转换异常", e);
}
}
/**
* 线程池状态
*/
@Data
@Accessors(chain = true)
public static class ThreadPoolStatus {
/**
* 核心线程数
*/
private int corePoolSize;
/**
* 最大线程数
*/
private int maxPoolSize;
/**
* 线程使用数
*/
private int activeThreads;
/**
* 线程空闲时间(秒)
*/
private int keepAliveSeconds;
/**
* 队列最大容量
*/
private int queueMaxCapacity;
/**
* 队列使用容量
*/
private int queueUseCapacity;
/**
* 队列剩余
*/
private int queueRemainingCapacity;
}
/**
* 用于解决配置转换问题
*/
@Retention(RetentionPolicy.RUNTIME)
@Target({ElementType.FIELD})
private @interface ConfigChangeEventToBean {
String value();
}
/**
* 动态线程池配置
* todo 需要自己改一下 注解的 value 属性
*/
@Data
@Accessors(chain = true)
static class DynamicPoolConfig {
/**
* 核心线程数
*/
@ConfigChangeEventToBean("dynamic.thread-pool.core-size")
private String coreSize;
/**
* 最大线程数
*/
@ConfigChangeEventToBean("dynamic.thread-pool.max-size")
private String maxSize;
/**
* 队列容量
*/
@ConfigChangeEventToBean("dynamic.thread-pool.queue-capacity")
private String queueCapacity;
/**
* 超时时间(秒)
*/
@ConfigChangeEventToBean("dynamic.thread-pool.keep-alive-seconds")
private String keepAliveSeconds;
}
/**
* 自定义线程池 实现自定义可变队列
*/
public static class ThreadPoolKafkaTaskExecutor extends ThreadPoolTaskExecutor {
public > void execute(T task) {
super.execute(task);
}
public > Future> submit(T task) {
return super.submit(task);
}
@Override
@NonNull
protected BlockingQueue createQueue(int queueCapacity) {
return queueCapacity > 0 ? new VariableLinkedBlockingQueue<>(queueCapacity) : new VariableLinkedBlockingQueue<>();
}
}
}
```
## 任务执行:UserModifyNameHandler.java
> 这个方法就是用于执行异步任务(不能是匿名内部类,不然反射获取不到对应的类就执行 kafka 的任务了)
>
> User.java 就是有两个成员变量 name 和 age 用于测试,这里就不展示了,可以看我 gitee 中的代码(链接再最上面的前言部分)
```java
@Slf4j
public class UserModifyNameHandler extends KafkaRejectPolicy.KafkaTaskInfoRunnable{
public UserModifyNameHandler(User taskData) {
super(taskData);
}
@Override
public void run() {
User taskData = super.getTaskData();
taskData.setName(taskData.getName() + "--修改了名字");
// 睡两秒钟模仿耗时
try {
Thread.sleep(2_000);
} catch (InterruptedException ignored) {
}
log.info("用户修改名字 任务开始执行 {}", taskData);
}
}
```
## 测试案例:KafkaPoolTestController.java
> 这个测试案例里面,我们只需要调用异常 /kafkaPoolTest/addTask 方法
>
> 会执行 10 条任务,我们再任务 UserModifyNameHandler 里面睡了两秒钟,
>
> 所以我们先设置的线程池参数(核心线程数:1,最大线程数:2,最大队列容量:2),只会执行 4 个任务
>
> 其他的就会被拒绝掉,然后放进 kafka 的队列里面,由 kafka 的消费者消费(KafkaConsumerService)
>
> 得到的效果就是线程池执行了 4 条任务(异步), kafka 执行了 6 条任务(同步)
```java
@RestController
@Slf4j
@RequestMapping("/kafkaPoolTest")
public class KafkaPoolTestController {
@Autowired
@Qualifier("dynamicKafkaRejectThreadPool")
private DynamicKafkaRejectThreadPoolConfig.ThreadPoolKafkaTaskExecutor dynamicThreadPool;
@Autowired
private DynamicThreadPoolConfig dynamicThreadPoolConfig;
/**
* 记录生成了多少用户
*/
private static final AtomicLong userCount = new AtomicLong(0);
/**
* 添加任务 测试 就让他等着 看拒绝情况
*/
@GetMapping("/addTask")
public Results addTask() {
// 一次性生成十个数据
List users = mockUsers();
// 线程池参数 核心线程:1 最大线程:2 队列容量:2
// 那么一次调用只能有 4 个任务执行,其他的都会被拒绝
// 那么拒绝的任务都进入了 kafka 的队列,由 KafkaConsumerService 处理
users.forEach(user -> {
Future> submit = dynamicThreadPool.submit(new UserModifyNameHandler(user));
});
return Results.success();
}
/**
* 获取参数
*/
@GetMapping("/getParameter")
public Results getParameter() {
DynamicThreadPoolConfig.ThreadPoolStatus threadPoolStatus = dynamicThreadPoolConfig.getThreadPoolStatus();
return Results.success(threadPoolStatus);
}
/**
* 模拟生成用户
*/
public synchronized static List mockUsers() {
List<@Nullable User> list = Lists.newArrayList();
for (int i = 0; i < 10; i++) {
long count = userCount.incrementAndGet();
User user = new User().setName("用户:" + count).setAge(18);
list.add(user);
}
return list;
}
}
```
## 测试结果:
> 由于打印的消息有点多,所以我就截取一部分进行展示
>
> 分别是
>
> 线程池异步进行了 4 条任务,用户 4、1、3、2
>
> kafka 同步执行了 6 条任务,用户 5、6、7、8、9、10(这里有可能也不是这个顺序执行的,这要看 kafka 入队顺序了,因为本来线程池就是异步的执行任务,kafka 为什么我要写成同步呢?因为本来进入 kafka 队列就是因为服务器资源紧张,线程池不够用了,再异步执行的话,跟线程池设置超多的最大线程数或者最大队列数有什么区别呢,这样会导致 OOM 的,所以这里使用同步,只让一个线程执行任务,就会减少服务器的压力,从而做到消峰+任务不丢失)
