@ -15,22 +15,22 @@
* limitations under the License .
* limitations under the License .
* /
* /
package cn.hippo4j. message.service ;
package cn.hippo4j. threadpool.alarm.handler ;
import cn.hippo4j.core.api.ThreadPoolCheckAlarm ;
import cn.hippo4j.common.executor.ThreadPoolExecutorHolder ;
import cn.hippo4j.common.executor.ThreadPoolExecutorRegistry ;
import cn.hippo4j.common.toolkit.CalculateUtil ;
import cn.hippo4j.common.toolkit.CalculateUtil ;
import cn.hippo4j.common.toolkit.ReflectUtil ;
import cn.hippo4j.common.toolkit.StringUtil ;
import cn.hippo4j.common.toolkit.StringUtil ;
import cn.hippo4j.core.executor.DynamicThreadPoolExecutor ;
import cn.hippo4j.threadpool.alarm.api.ThreadPoolCheckAlarm ;
import cn.hippo4j.core.executor.DynamicThreadPoolWrapper ;
import cn.hippo4j.threadpool.alarm.toolkit.ExecutorTraceContextUtil ;
import cn.hippo4j.core.executor.manage.GlobalThreadPoolManage ;
import cn.hippo4j.threadpool.message.core.request.AlarmNotifyRequest ;
import cn.hippo4j.core.executor.support.ThreadPoolBuilder ;
import cn.hippo4j.threadpool.message.api.NotifyTypeEnum ;
import cn.hippo4j.core.toolkit.ExecutorTraceContextUtil ;
import cn.hippo4j.threadpool.message.core.service.GlobalNotifyAlarmManage ;
import cn.hippo4j.core.toolkit.IdentifyUtil ;
import cn.hippo4j.threadpool.message.core.service.ThreadPoolNotifyAlarm ;
import cn.hippo4j.message.enums.NotifyTypeEnum ;
import cn.hippo4j.threadpool.message.core.service.ThreadPoolSendMessageService ;
import cn.hippo4j.message.request.AlarmNotifyRequest ;
import lombok.RequiredArgsConstructor ;
import lombok.RequiredArgsConstructor ;
import lombok.extern.slf4j.Slf4j ;
import lombok.extern.slf4j.Slf4j ;
import org.springframework.beans.factory.annotation.Value ;
import java.util.List ;
import java.util.List ;
import java.util.Objects ;
import java.util.Objects ;
@ -40,8 +40,16 @@ import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RejectedExecutionHandler ;
import java.util.concurrent.RejectedExecutionHandler ;
import java.util.concurrent.ScheduledExecutorService ;
import java.util.concurrent.ScheduledExecutorService ;
import java.util.concurrent.ScheduledThreadPoolExecutor ;
import java.util.concurrent.ScheduledThreadPoolExecutor ;
import java.util.concurrent.ThreadFactory ;
import java.util.concurrent.ThreadPoolExecutor ;
import java.util.concurrent.ThreadPoolExecutor ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.atomic.AtomicInteger ;
import static cn.hippo4j.common.propertie.EnvironmentProperties.active ;
import static cn.hippo4j.common.propertie.EnvironmentProperties.applicationName ;
import static cn.hippo4j.common.propertie.EnvironmentProperties.checkStateInterval ;
import static cn.hippo4j.common.propertie.EnvironmentProperties.itemId ;
import static cn.hippo4j.common.propertie.IdentifyProperties.IDENTIFY ;
/ * *
/ * *
* Default thread - pool check alarm handler .
* Default thread - pool check alarm handler .
@ -50,46 +58,42 @@ import java.util.concurrent.TimeUnit;
@RequiredArgsConstructor
@RequiredArgsConstructor
public class DefaultThreadPoolCheckAlarmHandler implements Runnable , ThreadPoolCheckAlarm {
public class DefaultThreadPoolCheckAlarmHandler implements Runnable , ThreadPoolCheckAlarm {
private final Hippo4jSendMessageService hippo4jSendMessageService ;
private final ThreadPoolSendMessageService threadPoolSendMessageService ;
@Value ( "${spring.profiles.active:UNKNOWN}" )
private String active ;
@Value ( "${spring.dynamic.thread-pool.item-id:}" )
private String itemId ;
@Value ( "${spring.application.name:UNKNOWN}" )
private String applicationName ;
@Value ( "${spring.dynamic.thread-pool.check-state-interval:5}" )
private Integer checkStateInterval ;
private final ScheduledExecutorService alarmNotifyExecutor = new ScheduledThreadPoolExecutor (
private final ScheduledExecutorService alarmNotifyExecutor = new ScheduledThreadPoolExecutor (
1 ,
1 ,
r - > new Thread ( r , "client.alarm.notify" ) ) ;
r - > new Thread ( r , "client.alarm.notify" ) ) ;
private final ExecutorService asyncAlarmNotifyExecutor = ThreadPoolBuilder . builder ( )
private final ExecutorService asyncAlarmNotifyExecutor = new ThreadPoolExecutor (
. poolThreadSize ( 2 , 4 )
2 ,
. threadFactory ( "client.execute.timeout.alarm" )
4 ,
. allowCoreThreadTimeOut ( true )
60L ,
. keepAliveTime ( 60L , TimeUnit . SECONDS )
TimeUnit . SECONDS ,
. workQueue ( new LinkedBlockingQueue ( 4096 ) )
new LinkedBlockingQueue < > ( 4096 ) ,
. rejected ( new ThreadPoolExecutor . AbortPolicy ( ) )
new ThreadFactory ( ) {
. build ( ) ;
private final AtomicInteger count = new AtomicInteger ( ) ;
@Override
@Override
public void run ( String . . . args ) throws Exception {
public Thread newThread ( Runnable r ) {
return new Thread ( "client.execute.timeout.alarm_" + count . incrementAndGet ( ) ) ;
}
} ,
new ThreadPoolExecutor . AbortPolicy ( ) ) ;
@Override
public void scheduleExecute ( ) {
alarmNotifyExecutor . scheduleWithFixedDelay ( this , 0 , checkStateInterval , TimeUnit . SECONDS ) ;
alarmNotifyExecutor . scheduleWithFixedDelay ( this , 0 , checkStateInterval , TimeUnit . SECONDS ) ;
}
}
@Override
@Override
public void run ( ) {
public void run ( ) {
List < String > listThreadPoolId = GlobalThreadPoolManage . listThreadPoolId ( ) ;
List < String > listThreadPoolId = ThreadPoolExecutorRegistry. listThreadPoolExecutor Id( ) ;
listThreadPoolId . forEach ( threadPoolId - > {
listThreadPoolId . forEach ( threadPoolId - > {
ThreadPoolNotifyAlarm threadPoolNotifyAlarm = GlobalNotifyAlarmManage . get ( threadPoolId ) ;
ThreadPoolNotifyAlarm threadPoolNotifyAlarm = GlobalNotifyAlarmManage . get ( threadPoolId ) ;
if ( threadPoolNotifyAlarm ! = null & & threadPoolNotifyAlarm . getAlarm ( ) ) {
if ( threadPoolNotifyAlarm ! = null & & threadPoolNotifyAlarm . getAlarm ( ) ) {
DynamicThreadPoolWrapper wrapper = GlobalThreadPoolManage . getExecutorService ( threadPoolId ) ;
ThreadPoolExecutorHolder executorHolder = ThreadPoolExecutorRegistry . getHolder ( threadPoolId ) ;
ThreadPoolExecutor executor = wrapp er. getExecutor ( ) ;
ThreadPoolExecutor executor = executorHold er. getExecutor ( ) ;
checkPoolCapacityAlarm ( threadPoolId , executor ) ;
checkPoolCapacityAlarm ( threadPoolId , executor ) ;
checkPoolActivityAlarm ( threadPoolId , executor ) ;
checkPoolActivityAlarm ( threadPoolId , executor ) ;
}
}
@ -116,7 +120,7 @@ public class DefaultThreadPoolCheckAlarmHandler implements Runnable, ThreadPoolC
if ( isSend ) {
if ( isSend ) {
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
hippo4j SendMessageService. sendAlarmMessage ( NotifyTypeEnum . CAPACITY , alarmNotifyRequest ) ;
threadPool SendMessageService. sendAlarmMessage ( NotifyTypeEnum . CAPACITY , alarmNotifyRequest ) ;
}
}
}
}
@ -139,7 +143,7 @@ public class DefaultThreadPoolCheckAlarmHandler implements Runnable, ThreadPoolC
if ( isSend ) {
if ( isSend ) {
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
hippo4j SendMessageService. sendAlarmMessage ( NotifyTypeEnum . ACTIVITY , alarmNotifyRequest ) ;
threadPool SendMessageService. sendAlarmMessage ( NotifyTypeEnum . ACTIVITY , alarmNotifyRequest ) ;
}
}
}
}
@ -155,11 +159,11 @@ public class DefaultThreadPoolCheckAlarmHandler implements Runnable, ThreadPoolC
if ( Objects . isNull ( alarmConfig ) | | ! alarmConfig . getAlarm ( ) ) {
if ( Objects . isNull ( alarmConfig ) | | ! alarmConfig . getAlarm ( ) ) {
return ;
return ;
}
}
ThreadPoolExecutor threadPoolExecutor = GlobalThreadPoolManage. getExecutorService ( threadPoolId ) . getExecutor ( ) ;
ThreadPoolExecutor threadPoolExecutor = ThreadPoolExecutorRegistry. getHolder ( threadPoolId ) . getExecutor ( ) ;
if ( threadPoolExecutor instanceof DynamicThreadPoolExecutor ) {
if ( Objects. equals ( threadPoolExecutor . getClass ( ) . getName ( ) , "cn.hippo4j.core.executor.DynamicThreadPoolExecutor" ) ) {
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
hippo4j SendMessageService. sendAlarmMessage ( NotifyTypeEnum . REJECT , alarmNotifyRequest ) ;
threadPool SendMessageService. sendAlarmMessage ( NotifyTypeEnum . REJECT , alarmNotifyRequest ) ;
}
}
} ;
} ;
asyncAlarmNotifyExecutor . execute ( checkPoolRejectedAlarmTask ) ;
asyncAlarmNotifyExecutor . execute ( checkPoolRejectedAlarmTask ) ;
@ -179,7 +183,6 @@ public class DefaultThreadPoolCheckAlarmHandler implements Runnable, ThreadPoolC
if ( Objects . isNull ( alarmConfig ) | | ! alarmConfig . getAlarm ( ) ) {
if ( Objects . isNull ( alarmConfig ) | | ! alarmConfig . getAlarm ( ) ) {
return ;
return ;
}
}
if ( threadPoolExecutor instanceof DynamicThreadPoolExecutor ) {
try {
try {
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
AlarmNotifyRequest alarmNotifyRequest = buildAlarmNotifyRequest ( threadPoolExecutor ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
alarmNotifyRequest . setThreadPoolId ( threadPoolId ) ;
@ -189,13 +192,12 @@ public class DefaultThreadPoolCheckAlarmHandler implements Runnable, ThreadPoolC
if ( StringUtil . isNotBlank ( executeTimeoutTrace ) ) {
if ( StringUtil . isNotBlank ( executeTimeoutTrace ) ) {
alarmNotifyRequest . setExecuteTimeoutTrace ( executeTimeoutTrace ) ;
alarmNotifyRequest . setExecuteTimeoutTrace ( executeTimeoutTrace ) ;
}
}
Runnable task = ( ) - > hippo4j SendMessageService. sendAlarmMessage ( NotifyTypeEnum . TIMEOUT , alarmNotifyRequest ) ;
Runnable task = ( ) - > threadPool SendMessageService. sendAlarmMessage ( NotifyTypeEnum . TIMEOUT , alarmNotifyRequest ) ;
asyncAlarmNotifyExecutor . execute ( task ) ;
asyncAlarmNotifyExecutor . execute ( task ) ;
} catch ( Throwable ex ) {
} catch ( Throwable ex ) {
log . error ( "Send thread pool execution timeout alarm error." , ex ) ;
log . error ( "Send thread pool execution timeout alarm error." , ex ) ;
}
}
}
}
}
/ * *
/ * *
* Build alarm notify request .
* Build alarm notify request .
@ -206,13 +208,17 @@ public class DefaultThreadPoolCheckAlarmHandler implements Runnable, ThreadPoolC
public AlarmNotifyRequest buildAlarmNotifyRequest ( ThreadPoolExecutor threadPoolExecutor ) {
public AlarmNotifyRequest buildAlarmNotifyRequest ( ThreadPoolExecutor threadPoolExecutor ) {
BlockingQueue < Runnable > blockingQueue = threadPoolExecutor . getQueue ( ) ;
BlockingQueue < Runnable > blockingQueue = threadPoolExecutor . getQueue ( ) ;
RejectedExecutionHandler rejectedExecutionHandler = threadPoolExecutor . getRejectedExecutionHandler ( ) ;
RejectedExecutionHandler rejectedExecutionHandler = threadPoolExecutor . getRejectedExecutionHandler ( ) ;
long rejectCount = threadPoolExecutor instanceof DynamicThreadPoolExecutor
long rejectCount = - 1L ;
? ( ( DynamicThreadPoolExecutor ) threadPoolExecutor ) . getRejectCountNum ( )
if ( Objects . equals ( threadPoolExecutor . getClass ( ) . getName ( ) , "cn.hippo4j.core.executor.DynamicThreadPoolExecutor" ) ) {
: - 1L ;
Object actualRejectCountNum = ReflectUtil . invoke ( threadPoolExecutor , "getRejectCountNum" ) ;
if ( actualRejectCountNum ! = null ) {
rejectCount = ( long ) actualRejectCountNum ;
}
}
return AlarmNotifyRequest . builder ( )
return AlarmNotifyRequest . builder ( )
. appName ( StringUtil . isBlank ( itemId ) ? applicationName : itemId )
. appName ( StringUtil . isBlank ( itemId ) ? applicationName : itemId )
. active ( active . toUpperCase ( ) )
. active ( active . toUpperCase ( ) )
. identify ( IdentifyUtil . getIdentify ( ) )
. identify ( I DENTIFY )
. corePoolSize ( threadPoolExecutor . getCorePoolSize ( ) )
. corePoolSize ( threadPoolExecutor . getCorePoolSize ( ) )
. maximumPoolSize ( threadPoolExecutor . getMaximumPoolSize ( ) )
. maximumPoolSize ( threadPoolExecutor . getMaximumPoolSize ( ) )
. poolSize ( threadPoolExecutor . getPoolSize ( ) )
. poolSize ( threadPoolExecutor . getPoolSize ( ) )