欢迎光临
我们一直在努力

深度解析 | DynamicTP 源码解析系列(七):如何管理第三方开源项目的线程池

前言

DynamicTP作为主流的动态线程池框架,核心能力是对线程池进行全生命周期管控、动态调参与监控告警,不仅能管理项目自定义线程池,更能适配各类第三方开源项目的线程池(Dubbo、Okhttp3 等),解决第三方线程池“难监控、难调参、易失控”的痛点。本文就聚焦DynamicTP如何实现对第三方开源项目线程池的管理。

一、通用实操流程(以okhttp3 为例)

无论适配哪种第三方开源项目,DynamicTP的管控落地都可遵循以下通用流程,兼顾简洁性与可扩展性,适配Spring Boot项目主流场景:

1.1 引入下述依赖

<!– SpringBoot1x、2x 用此依赖 –>
<dependency>
<groupId>org.dromara.dynamictp</groupId>
<artifactId>dynamic-tp-spring-boot-starter-adapter-okhttp3</artifactId>
<version>1.2.2</version>
</dependency>
<!– SpringBoot3x 用此依赖 –>
<dependency>
<groupId>org.dromara.dynamictp</groupId>
<artifactId>dynamic-tp-spring-boot-starter-adapter-okhttp3</artifactId>
<version>1.2.2-x</version>
</dependency>

1.2 配置文件中配置 okhttp3 线程池

dynamictp:
enabledCollect: true # 是否开启监控指标采集,默认false
collectorTypes: micrometer # 监控数据采集器类型(logging | micrometer | internal_logging | JMX),默认micrometer
monitorInterval: 5 # 监控时间间隔(报警判断、指标采集),默认5s
okhttp3Tp: # okhttp3 线程池配置
threadPoolName: okHttpClientTp
corePoolSize: 100
maximumPoolSize: 200
keepAliveTime: 60
runTimeout: 200
queueTimeout: 100
platformIds: [1,2] # 通知报警平台id,不配置默认拿上层platforms配置的所有平台
notifyItems: # 报警项,不配置自动会按默认值配置(变更通知、容量报警、活性报警、拒绝报警、任务超时报警)
type: change # 线程池核心参数变更通知
silencePeriod: 120 # 通知静默时间(单位:s),默认值1,0表示不静默

type: capacity # 队列容量使用率,报警项类型,查看源码 NotifyTypeEnum枚举类
threshold: 80 # 报警阈值,意思是队列使用率达到70%告警;默认值=70
count: 2 # 在一个统计周期内,如果触发阈值的数量达到 count,则触发报警;默认值=1
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默,默认值=120

type: liveness # 线程池活性
threshold: 80 # 报警阈值,意思是活性达到70%告警;默认值=70
count: 3 # 在一个统计周期内,如果触发阈值的数量达到 count,则触发报警;默认值=1
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默;默认值=120

type: reject # 触发任务拒绝告警
count: 1 # 在一个统计周期内,如果触发拒绝策略次数达到 count,则触发报警;默认值=1
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默;默认值=120

type: run_timeout # 任务执行超时告警
count: 20 # 在一个统计周期内,如果执行超时次数达到 count,则触发报警;默认值=10
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 30 # 报警静默时间(单位:s),0表示不静默;默认值=120

type: queue_timeout # 任务排队超时告警
count: 5 # 在一个统计周期内,如果排队超时次数达到 count,则触发报警;默认值=10
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默;默认值=120

1.3 启动日志

DynamicTp adapter, okhttp3 executors init end, executors: {okHttpClientTp=ExecutorWrapper(threadPoolName=okHttpClientTp, executor=java.util.concurrent.ThreadPoolExecutor@f336fd[Running, pool size = 0, active threads = 0, queued tasks = 0, completed tasks = 0], threadPoolAliasName=null, notifyItems=[NotifyItem(platforms=null, enabled=true, type=liveness, threshold=70, interval=120, clusterLimit=1), NotifyItem(platforms=null, enabled=true, type=change, threshold=0, interval=1, clusterLimit=1), NotifyItem(platforms=null, enabled=true, type=capacity, threshold=70, interval=120, clusterLimit=1)], notifyEnabled=true)}
DynamicTp okhttp3Tp adapter, [okHttpClientTp] refreshed end, changed keys: [corePoolSize, maxPoolSize], corePoolSize: [0 => 100], maxPoolSize: [2147483647 => 200], keepAliveTime: [60 => 60]

二、源码分析

2.1 Okhttp3TpAutoConfiguration

@Configuration
//ConditionalOnClass:基于「类路径是否存在指定类」的条件注解
@ConditionalOnClass(name = "okhttp3.OkHttpClient")
//ConditionalOnBean:基于「容器中是否存在指定 Bean」的条件注解
@ConditionalOnBean({DtpBaseBeanConfiguration.class})
//AutoConfigureAfter:控制「自动配置类的加载顺序」.指定当前自动配置类必须在指定的配置类加载完成后才加载
@AutoConfigureAfter({DtpBaseBeanConfiguration.class})
public class Okhttp3TpAutoConfiguration {

@Bean
@ConditionalOnMissingBean
public Okhttp3DtpAdapter okhttp3DtpAdapter() {
return new Okhttp3DtpAdapter();
}
}

2.2 Okhttp3DtpAdapter

public class Okhttp3DtpAdapter extends AbstractDtpAdapter {
//三方框架名称
private static final String TP_PREFIX = "okhttp3Tp";
//三方框架定义的线程池的字段名称
private static final String EXECUTOR_SERVICE_FIELD = "executorService";

private static final String EXECUTOR_SERVICE_FIELD_ALTERNATIVE = "executorServiceOrNull";

// 刷新配置信息
@Override
public void refresh(DtpProperties dtpProperties) {
//获取 Okhttp3 线程池的配置信息 (如果扩展其他三方组件,需要在 DtpProperties 添加对应的配置信息)
refresh(dtpProperties.getOkhttp3Tp(), dtpProperties.getPlatforms());
}

@Override
protected String getTpPrefix() {
return TP_PREFIX;
}

// 父类 AbstractDtpAdapter 的onContextRefreshedEvent 方法调用
@Override
protected void initialize() {
super.initialize();
// 从 spring容器 获取OkHttpClient 实例
val beans = ContextManagerHelper.getBeansOfType(OkHttpClient.class);
if (MapUtils.isEmpty(beans)) {
log.warn("Cannot find beans of type OkHttpClient.");
return;
}
//遍历 OkHttpClient 实例 , 获取线程池,通过装饰器模式增强,再通过 反射 替换
beans.forEach((k, v) -> {
val dispatcher = v.dispatcher();
//获取 okhttp3 中的线程池
val executor = dispatcher.executorService();
if (!(executor instanceof ThreadPoolExecutor)) {
return;
}

Field field = FieldUtils.getField(dispatcher.getClass(), EXECUTOR_SERVICE_FIELD, true);
if (Objects.isNull(field)) {
field = ReflectionUtil.getField(dispatcher.getClass(), EXECUTOR_SERVICE_FIELD_ALTERNATIVE);
}
// 将 okhttp3 中的线程池的线程池增强,再通过反射机制,修改 EXECUTOR_SERVICE_FIELD
enhanceOriginExecutor(genTpName(k), (ThreadPoolExecutor) executor, field, dispatcher);
});
}

private String genTpName(String clientName) {
return TP_PREFIX + "#" + clientName;
}
}

AbstractDtpAdapter

protected final Map<String, ExecutorWrapper> executors = Maps.newHashMap();
//注册监听器
protected AbstractDtpAdapter() {
EventBusManager.register(this);
}
//监听CustomContextRefreshedEvent 事件
@Subscribe
public synchronized void onContextRefreshedEvent(CustomContextRefreshedEvent event) {
try {
//从 spring容器 获取 DtpProperties 实例
DtpProperties dtpProperties = ContextManagerHelper.getBean(DtpProperties.class);
//初始化
initialize();
//初始化之后执行
afterInitialize();
//刷新配置
refresh(dtpProperties);
log.info("DynamicTp adapter, {} init end, executors {}", getTpPrefix(), executors.keySet());
} catch (Throwable e) {
log.error("DynamicTp adapter, {} init failed.", getTpPrefix(), e);
}
}
// 增强 原始的线程池
protected void enhanceOriginExecutor(String tpName, ThreadPoolExecutor executor, Field field, Object targetObj) {
//封装成 ThreadPoolExecutorProxy
ThreadPoolExecutorProxy proxy = new ThreadPoolExecutorProxy(executor);
boolean r = ReflectionUtil.setFieldValue(field, targetObj, proxy);
if (r) {
putAndFinalize(tpName, executor, proxy);
}
}
protected void putAndFinalize(String tpName, ExecutorService origin, Executor targetForWrapper) {
//executors : 维护线程池
executors.put(tpName, new ExecutorWrapper(tpName, targetForWrapper));
//关闭线程池
shutdownOriginalExecutor(origin);
}
//初始化之后执行
protected void afterInitialize() {
getExecutorWrappers().forEach((k, v) -> AwareManager.register(v));
}
//刷新配置
public void refresh(List<TpExecutorProps> propsList, List<NotifyPlatform> platforms) {
val executorWrappers = getExecutorWrappers();
if (CollectionUtils.isEmpty(propsList) || MapUtils.isEmpty(executorWrappers)) {
return;
}

val tmpMap = StreamUtil.toMap(propsList, TpExecutorProps::getThreadPoolName);
executorWrappers.forEach((k, v) -> refresh(v, platforms, tmpMap.get(k)));
}
//修改配置
public void refresh(ExecutorWrapper executorWrapper, List<NotifyPlatform> platforms, TpExecutorProps props) {
if (Objects.isNull(props) || Objects.isNull(executorWrapper) || containsInvalidParams(props, log)) {
return;
}
TpMainFields oldFields = getTpMainFields(executorWrapper, props);
doRefresh(executorWrapper, platforms, props);
TpMainFields newFields = getTpMainFields(executorWrapper, props);
if (oldFields.equals(newFields)) {
log.debug("DynamicTp adapter, main properties of [{}] have not changed.",
executorWrapper.getThreadPoolName());
return;
}

List<FieldInfo> diffFields = EQUATOR.getDiffFields(oldFields, newFields);
List<String> diffKeys = diffFields.stream().map(FieldInfo::getFieldName).collect(toList());
NoticeManager.tryNoticeAsync(executorWrapper, oldFields, diffKeys);
log.info("DynamicTp adapter, [{}] refreshed end, changed keys: {}, corePoolSize: [{}], "
+ "maxPoolSize: [{}], keepAliveTime: [{}]",
executorWrapper.getThreadPoolName(), diffKeys,
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getCorePoolSize(), newFields.getCorePoolSize()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getMaxPoolSize(), newFields.getMaxPoolSize()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getKeepAliveTime(), newFields.getKeepAliveTime()));
}

三、总结

服务启动会自动从 Spring 容器中获取所有被 Spring 容器管理的 OkHttpClient 实例。通过事件监听机制,修改线程池的核心配置。

线程池名称规则:beanName + Tp(可以在启动日志找输出的线程池名称)。

okhttp3 线程池只在异步请求时生效,同步请求不会使用 okhttp3 线程池。

okhttp3 线程池享有动态调参、监控、通知告警完整的功能。

队列大小不能调整

四、结尾

本文介绍了okhttp3 线程池的使步骤以及通过源码分析了DynamicTP管理第三方开源项目线程池的流程。如果在工作中遇到了要自定义适配未支持的第三方线程池,希望本文能提供一个比较清晰的思路。

思考:为什么说okhttp3中线程池的队列大小不能调整?

赞(0)
未经允许不得转载:171主机测评 » 深度解析 | DynamicTP 源码解析系列(七):如何管理第三方开源项目的线程池
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址