最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

PowerJob的OmsLogHandler工作流程源碼解析

 更新時間:2023年12月25日 08:30:42   作者:codecraft  
這篇文章主要為大家介紹了PowerJob的OmsLogHandler工作流程源碼解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

本文主要研究一下PowerJob的OmsLogHandler

OmsLogHandler

tech/powerjob/worker/background/OmsLogHandler.java

@Slf4j
public class OmsLogHandler {
    private final String workerAddress;
    private final Transporter transporter;
    private final ServerDiscoveryService serverDiscoveryService;
    // 處理線程,需要通過線程池啟動
    public final Runnable logSubmitter = new LogSubmitter();
    // 上報鎖,只需要一個線程上報即可
    private final Lock reportLock = new ReentrantLock();
    // 生產(chǎn)者消費者模式,異步上傳日志
    private final BlockingQueue<InstanceLogContent> logQueue = Queues.newLinkedBlockingQueue(10240);
    // 每次上報攜帶的數(shù)據(jù)條數(shù)
    private static final int BATCH_SIZE = 20;
    // 本地囤積閾值
    private static final int REPORT_SIZE = 1024;
    public OmsLogHandler(String workerAddress, Transporter transporter, ServerDiscoveryService serverDiscoveryService) {
        this.workerAddress = workerAddress;
        this.transporter = transporter;
        this.serverDiscoveryService = serverDiscoveryService;
    }
    /**
     * 提交日志
     * @param instanceId 任務(wù)實例ID
     * @param logContent 日志內(nèi)容
     */
    public void submitLog(long instanceId, LogLevel logLevel, String logContent) {
        if (logQueue.size() > REPORT_SIZE) {
            // 線程的生命周期是個不可循環(huán)的過程,一個線程對象結(jié)束了不能再次start,只能一直創(chuàng)建和銷毀
            new Thread(logSubmitter).start();
        }
        InstanceLogContent tuple = new InstanceLogContent(instanceId, System.currentTimeMillis(), logLevel.getV(), logContent);
        boolean offerRet = logQueue.offer(tuple);
        if (!offerRet) {
            log.warn("[OmsLogHandler] [{}] submit log failed, maybe your log speed is too fast!", instanceId);
        }
    }
    //......
}
OmsLogHandler提供了submitLog方法,它先判斷l(xiāng)ogQueue大小是否超過REPORT_SIZE(1024),超過則通過異步線程執(zhí)行l(wèi)ogSubmitter;接著將內(nèi)容包裝為InstanceLogContent,然后放入到logQueue

LogSubmitter

private class LogSubmitter implements Runnable {
        @Override
        public void run() {
            boolean lockResult = reportLock.tryLock();
            if (!lockResult) {
                return;
            }
            try {
                final String currentServerAddress = serverDiscoveryService.getCurrentServerAddress();
                // 當前無可用 Server
                if (StringUtils.isEmpty(currentServerAddress)) {
                    if (!logQueue.isEmpty()) {
                        logQueue.clear();
                        log.warn("[OmsLogHandler] because there is no available server to report logs which leads to queue accumulation, oms discarded all logs.");
                    }
                    return;
                }
                List<InstanceLogContent> logs = Lists.newLinkedList();
                while (!logQueue.isEmpty()) {
                    try {
                        InstanceLogContent logContent = logQueue.poll(100, TimeUnit.MILLISECONDS);
                        logs.add(logContent);
                        if (logs.size() >= BATCH_SIZE) {
                            WorkerLogReportReq req = new WorkerLogReportReq(workerAddress, Lists.newLinkedList(logs));
                            // 不可靠請求,WEB日志不追求極致
                            TransportUtils.reportLogs(req, currentServerAddress, transporter);
                            logs.clear();
                        }
                    }catch (Exception ignore) {
                        break;
                    }
                }
                if (!logs.isEmpty()) {
                    WorkerLogReportReq req = new WorkerLogReportReq(workerAddress, logs);
                    TransportUtils.reportLogs(req, currentServerAddress, transporter);
                }
            }finally {
                reportLock.unlock();
            }
        }
    }
LogSubmitter實現(xiàn)了Runnable接口,其run方法先通過reportLock加鎖,成功才繼續(xù),它通過serverDiscoveryService.getCurrentServerAddress()獲取當前server的地址,若獲取不到則清空logQueue;否則while循環(huán),每次從logQueue拉取InstanceLogContent,放到linkedList,超過BATCH_SIZE(20)則創(chuàng)建WorkerLogReportReq,通過TransportUtils.reportLogs(req, currentServerAddress, transporter)上報,然后清空linkedList,跳出循環(huán)之后再上報剩下的日志,最后釋放鎖

reportLogs

tech/powerjob/worker/common/utils/TransportUtils.java

public static void reportLogs(WorkerLogReportReq req, String address, Transporter transporter) {
        final URL url = easyBuildUrl(ServerType.SERVER, S4W_PATH, S4W_HANDLER_REPORT_LOG, address);
        transporter.tell(url, req);
    }
    public static URL easyBuildUrl(ServerType serverType, String rootPath, String handlerPath, String address) {
        HandlerLocation handlerLocation = new HandlerLocation()
                .setRootPath(rootPath)
                .setMethodPath(handlerPath);
        return new URL()
                .setServerType(serverType)
                .setAddress(Address.fromIpv4(address))
                .setLocation(handlerLocation);
    }
reportLogs先通過easyBuildUrl構(gòu)建URL,再通過transporter.tell(url, req)發(fā)送請求,rootPath為server,handlerPath為reportLog

tell

AkkaTransporter

tech/powerjob/remote/akka/AkkaTransporter.java

public void tell(URL url, PowerSerializable request) {
        ActorSelection actorSelection = fetchActorSelection(url);
        actorSelection.tell(request, null);
    }
AkkaTransporter直接使用actorSelection發(fā)送請求

VertxTransporter

tech/powerjob/remote/http/vertx/VertxTransporter.java

public void tell(URL url, PowerSerializable request) {
        post(url, request, null);
    }

    private &lt;T&gt; CompletionStage&lt;T&gt; post(URL url, PowerSerializable request, Class&lt;T&gt; clz) {
        final String host = url.getAddress().getHost();
        final int port = url.getAddress().getPort();
        final String path = url.getLocation().toPath();
        RequestOptions requestOptions = new RequestOptions()
                .setMethod(HttpMethod.POST)
                .setHost(host)
                .setPort(port)
                .setURI(path);
        // 獲取遠程服務(wù)器的HTTP連接
        Future&lt;HttpClientRequest&gt; httpClientRequestFuture = httpClient.request(requestOptions);
        // 轉(zhuǎn)換 -&gt; 發(fā)送請求獲取響應(yīng)
        Future&lt;HttpClientResponse&gt; responseFuture = httpClientRequestFuture.compose(httpClientRequest -&gt;
            httpClientRequest
                .putHeader(HttpHeaderNames.CONTENT_TYPE, HttpHeaderValues.APPLICATION_JSON)
                .send(JsonObject.mapFrom(request).toBuffer())
        );
        return responseFuture.compose(httpClientResponse -&gt; {
            // throw exception
            final int statusCode = httpClientResponse.statusCode();
            if (statusCode != HttpResponseStatus.OK.code()) {
                // CompletableFuture.get() 時會傳遞拋出該異常
                throw new RemotingException(String.format("request [host:%s,port:%s,url:%s] failed, status: %d, msg: %s",
                       host, port, path, statusCode, httpClientResponse.statusMessage()
                        ));
            }

            return httpClientResponse.body().compose(x -&gt; {

                if (clz == null) {
                    return Future.succeededFuture(null);
                }

                if (clz.equals(String.class)) {
                    return Future.succeededFuture((T) x.toString());
                }

                return Future.succeededFuture(x.toJsonObject().mapTo(clz));
            });
        }).toCompletionStage();
    }
VertxTransporter則使用post方法通過vertx的httpClient發(fā)送請求

processWorkerLogReport

tech/powerjob/server/core/handler/AbWorkerRequestHandler.java

@Handler(path = S4W_HANDLER_REPORT_LOG, processType = ProcessType.NO_BLOCKING)
    public void processWorkerLogReport(WorkerLogReportReq req) {
        WorkerLogReportEvent event = new WorkerLogReportEvent()
                .setWorkerAddress(req.getWorkerAddress())
                .setLogNum(req.getInstanceLogContents().size());
        try {
            processWorkerLogReport0(req, event);
            event.setStatus(WorkerLogReportEvent.Status.SUCCESS);
        } catch (RejectedExecutionException re) {
            event.setStatus(WorkerLogReportEvent.Status.REJECTED);
        } catch (Throwable t) {
            event.setStatus(WorkerLogReportEvent.Status.EXCEPTION);
            log.warn("[WorkerRequestHandler] process worker report failed!", t);
        } finally {
            monitorService.monitor(event);
        }
    }
processWorkerLogReport通過processWorkerLogReport0進行處理,最后通過monitorService.monitor(event)上報監(jiān)控

processWorkerLogReport0

tech/powerjob/server/core/handler/WorkerRequestHandlerImpl.java

@Override
    protected void processWorkerLogReport0(WorkerLogReportReq req, WorkerLogReportEvent event) {
        // 這個效率應(yīng)該不會拉垮吧...也就是一些判斷 + Map#get 吧...
        instanceLogService.submitLogs(req.getWorkerAddress(), req.getInstanceLogContents());
    }
processWorkerLogReport0通過instanceLogService.submitLogs進行上報

submitLogs

tech/powerjob/server/core/instance/InstanceLogService.java

/**
     * 提交日志記錄,持久化到本地數(shù)據(jù)庫中
     * @param workerAddress 上報機器地址
     * @param logs 任務(wù)實例運行時日志
     */
    @Async(value = PJThreadPool.LOCAL_DB_POOL)
    public void submitLogs(String workerAddress, List<InstanceLogContent> logs) {

        List<LocalInstanceLogDO> logList = logs.stream().map(x -> {
            instanceId2LastReportTime.put(x.getInstanceId(), System.currentTimeMillis());

            LocalInstanceLogDO y = new LocalInstanceLogDO();
            BeanUtils.copyProperties(x, y);
            y.setWorkerAddress(workerAddress);
            return y;
        }).collect(Collectors.toList());

        try {
            CommonUtils.executeWithRetry0(() -> localInstanceLogRepository.saveAll(logList));
        }catch (Exception e) {
            log.warn("[InstanceLogService] persistent instance logs failed, these logs will be dropped: {}.", logs, e);
        }
    }
InstanceLogService通過PJThreadPool.LOCAL_DB_POOL線程池進行異步,它通過localInstanceLogRepository.saveAll(logList)保存到本地數(shù)據(jù)庫

monitor

tech/powerjob/server/monitor/PowerJobMonitorService.java

public void monitor(Event event) {
        monitors.forEach(m -> m.record(event));
    }
monitor方法遍歷monitors,挨個執(zhí)行record

LogMonitor

tech/powerjob/server/monitor/monitors/LogMonitor.java

public void record(Event event) {
        MDC.put(MDC_KEY_SERVER_ID, String.valueOf(serverInfo.getId()));
        LoggerFactory.getLogger(event.type()).info(event.message());
    }
LogMonitor的record方法通過日志打印event信息

小結(jié)

PowerJob的OmsLogHandler提供了submitLog方法,它先判斷l(xiāng)ogQueue大小是否超過REPORT_SIZE(1024),超過則通過異步線程執(zhí)行l(wèi)ogSubmitter;接著將內(nèi)容包裝為InstanceLogContent,然后放入到logQueue;logSubmitter主要是執(zhí)行reportLogs,它先通過easyBuildUrl構(gòu)建URL,再通過transporter.tell(url, req)發(fā)送請求,rootPath為server,handlerPath為reportLog;服務(wù)端的processWorkerLogReport通過processWorkerLogReport0進行處理(通過localInstanceLogRepository.saveAll(logList)保存到本地數(shù)據(jù)庫),最后通過monitorService.monitor(event)上報監(jiān)控。

以上就是PowerJob的OmsLogHandler的詳細內(nèi)容,更多關(guān)于PowerJob的OmsLogHandler的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • IDEA中Git版本回退的兩種實現(xiàn)方案

    IDEA中Git版本回退的兩種實現(xiàn)方案

    作為開發(fā)者,代碼版本回退是日常高頻操作,IntelliJ IDEA集成了強大的Git工具鏈,但面對reset和revert兩種核心回退方案,許多開發(fā)者仍存在選擇困惑,本文將解析Reset與Revert兩種方案的操作細節(jié)及避坑指南,需要的朋友可以參考下
    2025-03-03
  • java開發(fā)gui教程之jframe監(jiān)聽窗體大小變化事件和jframe創(chuàng)建窗體

    java開發(fā)gui教程之jframe監(jiān)聽窗體大小變化事件和jframe創(chuàng)建窗體

    這篇文章主要介紹了java開發(fā)gui教程中jframe監(jiān)聽窗體大小變化事件和jframe創(chuàng)建窗體的示例,需要的朋友可以參考下
    2014-03-03
  • Java中的分布式事務(wù)Seata詳解

    Java中的分布式事務(wù)Seata詳解

    這篇文章主要介紹了Java中的分布式事務(wù)Seata詳解,Seata 是一款開源的分布式事務(wù)解決方案,致力于提供高性能和簡單易用的分布式事務(wù)服務(wù),Seata 將為用戶提供了 AT、TCC、SAGA 和 XA 事務(wù)模式,為用戶打造一站式的分布式解決方案,需要的朋友可以參考下
    2023-08-08
  • SpringBoot中的手動提交事務(wù)

    SpringBoot中的手動提交事務(wù)

    在Spring框架中使用@Transactional注解通常管理事務(wù),但在多線程環(huán)境下此方法失效,本文討論了手動事務(wù)的必要性及其實現(xiàn)方式,探討了Spring的七種事務(wù)傳播行為和數(shù)據(jù)庫的四大特性與隔離級別,了解這些可以幫助開發(fā)者在無法使用聲明式事務(wù)時
    2024-09-09
  • java中redissonClient 分布式鎖的使用

    java中redissonClient 分布式鎖的使用

    在集群的情況下,用戶多次請求接口時,存入的內(nèi)容可能會導(dǎo)致重復(fù),這時候就可以使用分布式鎖來限制,本文就來介紹一下java中redissonClient 分布式鎖的使用,感興趣的可以了解一下
    2024-03-03
  • IDEA插件(BindED)之查看class文件的十六進制

    IDEA插件(BindED)之查看class文件的十六進制

    這篇文章主要介紹了IDEA插件(BindED)之查看class文件的十六進制,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2023-12-12
  • Java EasyExcel讀寫excel如何解決poi讀取大文件內(nèi)存溢出問題

    Java EasyExcel讀寫excel如何解決poi讀取大文件內(nèi)存溢出問題

    這篇文章主要介紹了Java EasyExcel讀寫excel如何解決poi讀取大文件內(nèi)存溢出問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-06-06
  • java實現(xiàn)oracle插入當前時間的方法

    java實現(xiàn)oracle插入當前時間的方法

    這篇文章主要介紹了java實現(xiàn)oracle插入當前時間的方法,以實例形式對比分析了java使用Oracle操作時間的技巧,具有一定參考借鑒價值,需要的朋友可以參考下
    2015-03-03
  • 詳解Java面試官最愛問的volatile關(guān)鍵字

    詳解Java面試官最愛問的volatile關(guān)鍵字

    這篇文章主要介紹了詳解Java面試官最愛問的volatile關(guān)鍵字,小編覺得還是挺不錯的,具有一定借鑒價值,需要的朋友可以參考下
    2018-01-01
  • 使用java實現(xiàn)猜拳小游戲

    使用java實現(xiàn)猜拳小游戲

    這篇文章主要為大家詳細介紹了使用java實現(xiàn)猜拳小游戲,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-07-07

最新評論

洞口县| 景东| 张北县| 巴青县| 观塘区| 荆州市| 高碑店市| 绥化市| 文成县| 平潭县| 拜城县| 常德市| 兴城市| 定远县| 铜鼓县| 永顺县| 昌江| 夹江县| 即墨市| 保德县| 深州市| 金华市| 扬中市| 台南县| 盖州市| 商河县| 兰州市| 平遥县| 定南县| 仙游县| 普安县| 炎陵县| 竹山县| 嵊州市| 通化市| 大连市| 曲阳县| 通江县| 贵港市| 静宁县| 威信县|