From 6c9a9f9f1102c238eef265f9b99eb460f57d8fb6 Mon Sep 17 00:00:00 2001 From: wtt <1136220284@qq.com> Date: Mon, 22 Apr 2024 14:40:44 +0800 Subject: [PATCH 1/2] refactor: add stream pull nacos config --- .../mone/log/stream/config/ConfigManager.java | 17 +++--- .../mone/log/stream/job/PullConfigJob.java | 59 +++++++++++++++++++ 2 files changed, 66 insertions(+), 10 deletions(-) create mode 100644 ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/job/PullConfigJob.java diff --git a/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java b/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java index dcbf88acc..60d9e779d 100644 --- a/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java +++ b/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java @@ -77,7 +77,7 @@ public class ConfigManager { * * @throws StreamException */ - public void initStream() throws StreamException { + public void initializeStreamConfig() throws StreamException { log.debug("[initStream} nacos dataId:{},group:{}", spaceDataId, DEFAULT_GROUP_ID); String streamConfigStr = nacosConfig.getConfigStr(spaceDataId, DEFAULT_GROUP_ID, DEFAULT_TIME_OUT_MS); MiLogStreamConfig milogStreamConfig; @@ -96,18 +96,15 @@ public void initStream() throws StreamException { for (Long spaceId : milogStreamDataMap.keySet()) { final String dataId = milogStreamDataMap.get(spaceId); // init spaceData config - String milogSpaceDataStr = nacosConfig.getConfigStr(dataId, DEFAULT_GROUP_ID, DEFAULT_TIME_OUT_MS); - if (StringUtils.isNotEmpty(milogSpaceDataStr)) { - MilogSpaceData milogSpaceData = GSON.fromJson(milogSpaceDataStr, MilogSpaceData.class); - if (null != milogSpaceData) { + String logSpaceDataStr = nacosConfig.getConfigStr(dataId, DEFAULT_GROUP_ID, DEFAULT_TIME_OUT_MS); + if (StringUtils.isNotEmpty(logSpaceDataStr)) { + MilogSpaceData milogSpaceData = GSON.fromJson(logSpaceDataStr, MilogSpaceData.class); + if (null != milogSpaceData && !milogSpaceDataMap.containsKey(spaceId)) { + MilogConfigListener configListener = new MilogConfigListener(spaceId, dataId, DEFAULT_GROUP_ID, milogSpaceData, nacosConfig); + addListener(spaceId, configListener); milogSpaceDataMap.put(spaceId, milogSpaceData); } } - MilogSpaceData milogSpaceData = milogSpaceDataMap.get(spaceId); - if (null != milogSpaceData) { - MilogConfigListener configListener = new MilogConfigListener(spaceId, dataId, DEFAULT_GROUP_ID, milogSpaceData, nacosConfig); - addListener(spaceId, configListener); - } log.info("[ConfigManager.initStream] added log config listener for spaceId:{},dataId:{}", spaceId, dataId); } } else { diff --git a/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/job/PullConfigJob.java b/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/job/PullConfigJob.java new file mode 100644 index 000000000..c1e70383a --- /dev/null +++ b/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/job/PullConfigJob.java @@ -0,0 +1,59 @@ +/* + * Copyright 2020 Xiaomi + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.xiaomi.mone.log.stream.job; + +import cn.hutool.core.thread.ThreadUtil; +import com.xiaomi.mone.log.stream.config.ConfigManager; +import com.xiaomi.youpin.docean.anno.Component; +import lombok.extern.slf4j.Slf4j; + +import javax.annotation.Resource; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +/** + * + * @description Regularly pull configuration comparison to compensate for the situation where nacos monitoring does not work. + * @version 1.0 + * @author wtt + * @date 2024/4/22 14:08 + * + */ +@Component +@Slf4j +public class PullConfigJob { + + @Resource + private ConfigManager configManager; + + public void init() { + + log.info("PullConfigJob execute"); + ScheduledExecutorService scheduledExecutor = Executors.newSingleThreadScheduledExecutor( + ThreadUtil.newNamedThreadFactory("pull-config", false) + ); + long initDelay = 2; + long intervalTime = 5; + scheduledExecutor.scheduleAtFixedRate(() -> { + try { + configManager.initializeStreamConfig(); + } catch (Exception e) { + log.error("PullConfigJob execute error", e); + } + }, initDelay, intervalTime, TimeUnit.MINUTES); + } +} From a954cc052703396e4bf2ac7767d98f24b5217056 Mon Sep 17 00:00:00 2001 From: wtt <1136220284@qq.com> Date: Mon, 22 Apr 2024 14:47:49 +0800 Subject: [PATCH 2/2] refactor: update log location --- .../java/com/xiaomi/mone/log/stream/config/ConfigManager.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java b/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java index 60d9e779d..110b99e39 100644 --- a/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java +++ b/ozhera-log/log-stream/src/main/java/com/xiaomi/mone/log/stream/config/ConfigManager.java @@ -103,9 +103,9 @@ public void initializeStreamConfig() throws StreamException { MilogConfigListener configListener = new MilogConfigListener(spaceId, dataId, DEFAULT_GROUP_ID, milogSpaceData, nacosConfig); addListener(spaceId, configListener); milogSpaceDataMap.put(spaceId, milogSpaceData); + log.info("[ConfigManager.initStream] added log config listener for spaceId:{},dataId:{}", spaceId, dataId); } } - log.info("[ConfigManager.initStream] added log config listener for spaceId:{},dataId:{}", spaceId, dataId); } } else { log.info("server start current not contain space config,uniqueMark:{}", uniqueMark);