本文目录导读:

我来详细介绍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<>();
}
}
这个实现包含了:
- 告警Appender:自定义logback Appender捕获日志
- 多渠道通知:邮件、钉钉、企业微信等
- 告警规则引擎:灵活的规则配置
- 告警聚合:减少告警风暴
- 告警降级:保护系统
- Spring Boot集成:简化配置
你可以根据实际需求选择合适的告警方式和规则进行配置。