Java日志告警案例如何实现

wen java案例 30

本文目录导读:

Java日志告警案例如何实现

  1. 基础日志告警框架
  2. 告警处理器实现
  3. 告警规则引擎
  4. 告警聚合与降级
  5. Spring Boot集成配置
  6. 使用示例

我来详细介绍Java日志告警的实现方案,包括多种实战案例。

基础日志告警框架

自定义日志Appender

import ch.qos.logback.core.AppenderBase;
import ch.qos.logback.classic.spi.ILoggingEvent;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
public class AlertAppender extends AppenderBase<ILoggingEvent> {
    private static final BlockingQueue<LogEvent> ALERT_QUEUE = 
        new LinkedBlockingQueue<>(1000);
    private AlertHandler alertHandler;
    private String thresholdLevel = "ERROR";
    @Override
    protected void append(ILoggingEvent event) {
        if (event.getLevel().levelStr.equalsIgnoreCase(thresholdLevel)) {
            LogEvent logEvent = new LogEvent(
                event.getFormattedMessage(),
                event.getLevel().levelStr,
                event.getLoggerName(),
                event.getTimeStamp(),
                Thread.currentThread().getName()
            );
            if (!ALERT_QUEUE.offer(logEvent)) {
                // 队列满时的降级处理
                logEvent.setMessage("[QUEUE_FULL] " + logEvent.getMessage());
            }
            // 异步处理告警
            alertHandler.handleAlert(logEvent);
        }
    }
    public void setThresholdLevel(String level) {
        this.thresholdLevel = level;
    }
    public void setAlertHandler(AlertHandler handler) {
        this.alertHandler = handler;
    }
}
class LogEvent {
    private String message;
    private String level;
    private String loggerName;
    private long timestamp;
    private String threadName;
    // 构造方法、getter/setter省略
}

logback.xml配置

<?xml version="1.0" encoding="UTF-8"?>
<configuration>
    <!-- 自定义告警Appender -->
    <appender name="ALERT" class="com.example.AlertAppender">
        <thresholdLevel>ERROR</thresholdLevel>
        <alertHandler>com.example.AlertHandler</alertHandler>
    </appender>
    <!-- 控制台Appender -->
    <appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
        <encoder>
            <pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
        </encoder>
    </appender>
    <!-- 文件Appender -->
    <appender name="FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
        <file>logs/app.log</file>
        <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
            <fileNamePattern>logs/app.%d{yyyy-MM-dd}.log</fileNamePattern>
            <maxHistory>30</maxHistory>
        </rollingPolicy>
        <encoder>
            <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
        </encoder>
    </appender>
    <root level="INFO">
        <appender-ref ref="CONSOLE"/>
        <appender-ref ref="FILE"/>
        <appender-ref ref="ALERT"/>
    </root>
</configuration>

告警处理器实现

多渠道告警处理器

public interface AlertHandler {
    void handleAlert(LogEvent event);
}
public class CompositeAlertHandler implements AlertHandler {
    private List<AlertHandler> handlers = new ArrayList<>();
    public void addHandler(AlertHandler handler) {
        handlers.add(handler);
    }
    @Override
    public void handleAlert(LogEvent event) {
        for (AlertHandler handler : handlers) {
            try {
                handler.handleAlert(event);
            } catch (Exception e) {
                // 单个处理器失败不影响其他
                System.err.println("Alert handler failed: " + e.getMessage());
            }
        }
    }
}

邮件告警处理器

import org.springframework.mail.SimpleMailMessage;
import org.springframework.mail.javamail.JavaMailSender;
public class EmailAlertHandler implements AlertHandler {
    private JavaMailSender mailSender;
    private String[] alertEmails;
    private String applicationName;
    @Override
    public void handleAlert(LogEvent event) {
        if (shouldSendAlert(event)) {
            sendAlertEmail(event);
        }
    }
    private boolean shouldSendAlert(LogEvent event) {
        // 1. 错误频率检查
        // 2. 错误类型过滤
        // 3. 时间窗口限制
        return true;
    }
    private void sendAlertEmail(LogEvent event) {
        SimpleMailMessage message = new SimpleMailMessage();
        message.setTo(alertEmails);
        message.setSubject(String.format("[%s] 告警: %s", 
            applicationName, event.getLevel()));
        message.setText(buildEmailContent(event));
        try {
            mailSender.send(message);
        } catch (Exception e) {
            System.err.println("Failed to send alert email: " + e.getMessage());
        }
    }
    private String buildEmailContent(LogEvent event) {
        StringBuilder sb = new StringBuilder();
        sb.append("=== 系统告警通知 ===\n");
        sb.append("应用名称: ").append(applicationName).append("\n");
        sb.append("告警级别: ").append(event.getLevel()).append("\n");
        sb.append("发生时间: ").append(new Date(event.getTimestamp())).append("\n");
        sb.append("线程名称: ").append(event.getThreadName()).append("\n");
        sb.append("记录器: ").append(event.getLoggerName()).append("\n");
        sb.append("告警内容: ").append(event.getMessage()).append("\n");
        sb.append("==================");
        return sb.toString();
    }
}

钉钉/企业微信机器人告警

import okhttp3.*;
import com.fasterxml.jackson.databind.ObjectMapper;
public class DingTalkAlertHandler implements AlertHandler {
    private String webhookUrl;
    private String secret;
    private OkHttpClient client = new OkHttpClient();
    private ObjectMapper mapper = new ObjectMapper();
    @Override
    public void handleAlert(LogEvent event) {
        try {
            String json = buildDingTalkMessage(event);
            RequestBody body = RequestBody.create(
                MediaType.parse("application/json; charset=utf-8"), 
                json
            );
            Request request = new Request.Builder()
                .url(webhookUrl)
                .post(body)
                .build();
            client.newCall(request).enqueue(new Callback() {
                @Override
                public void onFailure(Call call, IOException e) {
                    System.err.println("DingTalk alert failed: " + e.getMessage());
                }
                @Override
                public void onResponse(Call call, Response response) throws IOException {
                    if (!response.isSuccessful()) {
                        System.err.println("DingTalk alert response: " + response.body().string());
                    }
                }
            });
        } catch (Exception e) {
            System.err.println("Failed to send DingTalk alert: " + e.getMessage());
        }
    }
    private String buildDingTalkMessage(LogEvent event) throws Exception {
        Map<String, Object> message = new HashMap<>();
        message.put("msgtype", "markdown");
        Map<String, Object> markdown = new HashMap<>();
        markdown.put("title", "系统告警");
        markdown.put("text", String.format(
            "### 系统告警通知\n" +
            "- **告警级别**: %s\n" +
            "- **发生时间**: %s\n" +
            "- **错误内容**: %s\n",
            event.getLevel(),
            new Date(event.getTimestamp()).toString(),
            event.getMessage()
        ));
        message.put("markdown", markdown);
        return mapper.writeValueAsString(message);
    }
}

告警规则引擎

规则定义

public class AlertRule {
    private String name;
    private String pattern;  // 日志匹配模式
    private String level;     // 告警级别
    private int threshold;    // 触发阈值
    private long timeWindow;  // 时间窗口(秒)
    private List<String> channels; // 告警渠道
    // 构造方法、getter/setter
}

规则匹配引擎

import com.google.common.cache.Cache;
import com.google.common.cache.CacheBuilder;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
public class AlertRuleEngine {
    private List<AlertRule> rules;
    private Cache<String, AtomicInteger> errorCounters;
    public AlertRuleEngine() {
        this.rules = new ArrayList<>();
        this.errorCounters = CacheBuilder.newBuilder()
            .expireAfterWrite(1, TimeUnit.MINUTES)
            .build();
    }
    public void addRule(AlertRule rule) {
        rules.add(rule);
    }
    public boolean evaluate(LogEvent event) {
        for (AlertRule rule : rules) {
            if (matchesRule(event, rule)) {
                String key = generateKey(event, rule);
                AtomicInteger counter = errorCounters.getIfPresent(key);
                if (counter == null) {
                    counter = new AtomicInteger(0);
                    errorCounters.put(key, counter);
                }
                int currentCount = counter.incrementAndGet();
                if (currentCount >= rule.getThreshold()) {
                    // 重置计数器
                    counter.set(0);
                    return true; // 触发告警
                }
            }
        }
        return false;
    }
    private boolean matchesRule(LogEvent event, AlertRule rule) {
        // 1. 级别匹配
        if (!event.getLevel().equalsIgnoreCase(rule.getLevel())) {
            return false;
        }
        // 2. 模式匹配
        if (rule.getPattern() != null && 
            !event.getMessage().contains(rule.getPattern())) {
            return false;
        }
        return true;
    }
    private String generateKey(LogEvent event, AlertRule rule) {
        return rule.getName() + ":" + event.getLoggerName();
    }
}

告警聚合与降级

告警聚合器

import java.util.concurrent.*;
import java.util.stream.Collectors;
public class AlertAggregator {
    private final ScheduledExecutorService scheduler = 
        Executors.newScheduledThreadPool(1);
    private final BlockingQueue<LogEvent> alertQueue = 
        new LinkedBlockingQueue<>(10000);
    public void start(int aggregationIntervalSeconds) {
        scheduler.scheduleAtFixedRate(
            this::processAlerts,
            aggregationIntervalSeconds,
            aggregationIntervalSeconds,
            TimeUnit.SECONDS
        );
    }
    private void processAlerts() {
        List<LogEvent> currentBatch = new ArrayList<>();
        alertQueue.drainTo(currentBatch);
        if (currentBatch.isEmpty()) return;
        // 按错误类型分组
        Map<String, List<LogEvent>> grouped = currentBatch.stream()
            .collect(Collectors.groupingBy(LogEvent::getMessage));
        // 生成聚合告警
        for (Map.Entry<String, List<LogEvent>> entry : grouped.entrySet()) {
            if (entry.getValue().size() >= 5) { // 聚合阈值
                sendAggregatedAlert(entry);
            }
        }
    }
    private void sendAggregatedAlert(Map.Entry<String, List<LogEvent>> entry) {
        List<LogEvent> events = entry.getValue();
        LogEvent first = events.get(0);
        String aggregatedMessage = String.format(
            "[聚合告警] 类型: %s, 发生次数: %d, 时间范围: %s - %s",
            first.getMessage(),
            events.size(),
            new Date(events.get(0).getTimestamp()),
            new Date(events.get(events.size() - 1).getTimestamp())
        );
        // 发送聚合后的告警
        LogEvent aggregatedEvent = new LogEvent(
            aggregatedMessage,
            "ERROR",
            first.getLoggerName(),
            System.currentTimeMillis(),
            "aggregator"
        );
        // 调用告警处理器
    }
}

告警降级处理器

public class AlertDegradationHandler {
    private final AtomicInteger alertCounter = new AtomicInteger(0);
    private final AtomicLong lastAlertTime = new AtomicLong(0);
    private volatile boolean degraded = false;
    private static final int MAX_ALERTS_PER_MINUTE = 10;
    private static final long DEGRADATION_DURATION_MS = 60_000;
    public boolean shouldDegrade() {
        long currentTime = System.currentTimeMillis();
        // 重置计数器(如果超过1分钟)
        if (currentTime - lastAlertTime.get() > 60_000) {
            alertCounter.set(0);
            lastAlertTime.set(currentTime);
            degraded = false;
            return false;
        }
        int currentCount = alertCounter.incrementAndGet();
        if (currentCount > MAX_ALERTS_PER_MINUTE && !degraded) {
            degraded = true;
            return true;
        }
        return degraded;
    }
    public boolean isDegraded() {
        if (degraded) {
            // 检查降级是否结束
            if (System.currentTimeMillis() - lastAlertTime.get() > DEGRADATION_DURATION_MS) {
                degraded = false;
                alertCounter.set(0);
                return false;
            }
            return true;
        }
        return false;
    }
}

Spring Boot集成配置

自动配置类

import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;
@Configuration
@ConfigurationProperties(prefix = "alert")
public class AlertProperties {
    private boolean enabled = true;
    private String level = "ERROR";
    private List<String> channels = Arrays.asList("email");
    private List<String> emails;
    private String dingtalkWebhook;
    private int aggregationSeconds = 60;
    private int threshold = 5;
    // getter/setter
}
@Configuration
public class AlertAutoConfiguration {
    @Bean
    public AlertRuleEngine alertRuleEngine() {
        return new AlertRuleEngine();
    }
    @Bean
    public CompositeAlertHandler alertHandler(AlertProperties properties) {
        CompositeAlertHandler composite = new CompositeAlertHandler();
        if (properties.getChannels().contains("email")) {
            // 添加邮件处理器
        }
        if (properties.getChannels().contains("dingtalk")) {
            DingTalkAlertHandler dingtalk = new DingTalkAlertHandler();
            dingtalk.setWebhookUrl(properties.getDingtalkWebhook());
            composite.addHandler(dingtalk);
        }
        return composite;
    }
    @Bean
    public AlertAggregator alertAggregator(AlertProperties properties) {
        AlertAggregator aggregator = new AlertAggregator();
        aggregator.start(properties.getAggregationSeconds());
        return aggregator;
    }
}

application.yml配置

alert:
  enabled: true
  level: ERROR
  channels:
    - email
    - dingtalk
  emails:
    - admin1@example.com
    - admin2@example.com
  dingtalk-webhook: https://oapi.dingtalk.com/robot/send?access_token=xxx
  aggregation-seconds: 60
  threshold: 5
  # 规则配置
  rules:
    - name: database-error
      pattern: "DatabaseException"
      level: ERROR
      threshold: 3
      time-window: 60
      channels:
        - email
        - dingtalk
    - name: timeout-alert
      pattern: "TimeoutException"
      level: WARN
      threshold: 10
      time-window: 300
      channels:
        - dingtalk

使用示例

业务代码中使用

@Service
public class OrderService {
    private static final Logger logger = LoggerFactory.getLogger(OrderService.class);
    @Autowired
    private AlertRuleEngine ruleEngine;
    public void processOrder(Order order) {
        try {
            // 业务逻辑
        } catch (DatabaseException e) {
            logger.error("订单处理数据库异常: {}", e.getMessage(), e);
            // 如果需要立即告警,可以手动触发
            LogEvent event = new LogEvent(
                e.getMessage(),
                "ERROR",
                this.getClass().getName(),
                System.currentTimeMillis(),
                Thread.currentThread().getName()
            );
            if (ruleEngine.evaluate(event)) {
                // 触发告警(会自动通过Appender处理)
            }
        }
    }
}

监控告警配置

@RestController
@RequestMapping("/api/alerts")
public class AlertController {
    @Autowired
    private AlertRuleEngine ruleEngine;
    @GetMapping("/rules")
    public List<AlertRule> getRules() {
        // 返回当前告警规则
        return new ArrayList<>();
    }
    @PostMapping("/rules")
    public void addRule(@RequestBody AlertRule rule) {
        ruleEngine.addRule(rule);
    }
    @GetMapping("/stats")
    public Map<String, Object> getStats() {
        // 返回告警统计信息
        return new HashMap<>();
    }
}

这个实现包含了:

  1. 告警Appender:自定义logback Appender捕获日志
  2. 多渠道通知:邮件、钉钉、企业微信等
  3. 告警规则引擎:灵活的规则配置
  4. 告警聚合:减少告警风暴
  5. 告警降级:保护系统
  6. Spring Boot集成:简化配置

你可以根据实际需求选择合适的告警方式和规则进行配置。

抱歉,评论功能暂时关闭!