feat:更新日志存储和查询模块
This commit is contained in:
Binary file not shown.
Binary file not shown.
@@ -1,16 +1,19 @@
|
||||
from collections import defaultdict
|
||||
from typing import List
|
||||
|
||||
from sqlalchemy import insert, select
|
||||
from sqlalchemy import insert, select, func
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from models.log_alert import LogAlert
|
||||
from schemas.log_alert import LogAlertTrigger
|
||||
from schemas.log_alert import LogAlertTrigger, LogAlertResponse
|
||||
from service import openobserve_service
|
||||
from service.message_service import send_wechat_message, MessageEnum
|
||||
from utils.time import get_timestamp_range, format_timestamp
|
||||
from config.logging import logger
|
||||
|
||||
|
||||
def trigger_alert(log_alter: LogAlertTrigger, db: Session) -> bool:
|
||||
logger.info(f"{format_timestamp(log_alter.timestamp)} 发生告警:{log_alter.alter_name}")
|
||||
|
||||
start_time, end_time = get_timestamp_range(log_alter.timestamp, 60)
|
||||
result = openobserve_service.get_log(start_time, end_time)
|
||||
hits = result.get("hits")
|
||||
@@ -20,20 +23,19 @@ def trigger_alert(log_alter: LogAlertTrigger, db: Session) -> bool:
|
||||
return False
|
||||
|
||||
# 判断日志告警是否已经存在
|
||||
log_db_alter = db.execute(
|
||||
select(LogAlert)
|
||||
.where(LogAlert.alter_timestamp == log_alter.timestamp)
|
||||
).first()
|
||||
stmt = select(LogAlert).where(LogAlert.alter_time == format_timestamp(log_alter.timestamp))
|
||||
log_db_alter = db.execute(stmt).first()
|
||||
|
||||
if log_db_alter is not None:
|
||||
return False
|
||||
|
||||
# 企业微信通知
|
||||
alter_message = f"🔔 **告警服务**:{log_alter.alter_name} \n 🕒 **告警时间**:{format_timestamp(log_alter.timestamp)}"
|
||||
send_wechat_message(MessageEnum.MARKDOWN2, alter_message)
|
||||
|
||||
# 批量插入
|
||||
return batch_insert(log_alter, hits, db)
|
||||
batch_insert(log_alter, hits, db)
|
||||
|
||||
# 通知消息
|
||||
# notify_message(log_alter)
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool:
|
||||
@@ -41,10 +43,11 @@ def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool:
|
||||
[
|
||||
{
|
||||
"alter_name": log_alter.alter_name,
|
||||
"alter_timestamp": log_alter.timestamp,
|
||||
"timestamp": hit["_timestamp"],
|
||||
"message": hit["message"],
|
||||
"status": "pending"
|
||||
"alter_time": format_timestamp(log_alter.timestamp),
|
||||
"log_time": format_timestamp(hit["_timestamp"]),
|
||||
"log_topic": hit["topic"],
|
||||
"log_message": hit["message"],
|
||||
"log_status": "pending"
|
||||
}
|
||||
for hit in hits
|
||||
]
|
||||
@@ -56,15 +59,44 @@ def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def notify_message(log_alter: LogAlertTrigger):
|
||||
# 企业微信通知
|
||||
alter_message = f"🔔 **告警服务**:{log_alter.alter_name} \n 🕒 **告警时间**:{format_timestamp(log_alter.timestamp)}"
|
||||
send_wechat_message(MessageEnum.MARKDOWN2, alter_message)
|
||||
|
||||
|
||||
def query_alert(db: Session):
|
||||
alerts = db.execute(select(LogAlert)).scalars().all()
|
||||
stmt = (
|
||||
select(
|
||||
LogAlert.alter_name.label("alterName"),
|
||||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.alter_time).label("alterTime"),
|
||||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.log_time).label("logTime"),
|
||||
LogAlert.log_topic.label("logTopic"),
|
||||
LogAlert.log_message.label("logMessage"),
|
||||
LogAlert.log_status.label("logStatus")
|
||||
)
|
||||
.order_by(LogAlert.alter_name.desc(), LogAlert.alter_time.desc())
|
||||
)
|
||||
|
||||
# 按照name和timestamp分类
|
||||
classified = defaultdict(lambda: defaultdict(list))
|
||||
for alert in alerts:
|
||||
classified[alert.alter_name][alert.alter_timestamp].append({
|
||||
"timestamp": format_timestamp(alert.timestamp),
|
||||
"message": alert.message
|
||||
})
|
||||
results = db.execute(stmt).fetchall()
|
||||
|
||||
return dict(classified)
|
||||
return [LogAlertResponse.model_validate(result) for result in results]
|
||||
|
||||
|
||||
def query_alert_by_name(name: str, db: Session) -> List[LogAlertResponse]:
|
||||
stmt = (
|
||||
select(
|
||||
LogAlert.alter_name.label("alterName"),
|
||||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.alter_time).label("alterTime"),
|
||||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.log_time).label("logTime"),
|
||||
LogAlert.log_topic.label("logTopic"),
|
||||
LogAlert.log_message.label("logMessage"),
|
||||
LogAlert.log_status.label("logStatus")
|
||||
)
|
||||
.where(LogAlert.alter_name == name)
|
||||
.order_by(LogAlert.alter_time.desc(), LogAlert.log_time.asc())
|
||||
)
|
||||
|
||||
results = db.execute(stmt).fetchall()
|
||||
|
||||
return [LogAlertResponse.model_validate(result) for result in results]
|
||||
|
||||
@@ -26,7 +26,8 @@ def send_wechat_message(message_type: MessageEnum, message: str):
|
||||
settings.wechat_webhook_url,
|
||||
headers={"Content-Type": "application/json"},
|
||||
params={'key': settings.wechat_webhook_key},
|
||||
data=json.dumps(wechat_message, ensure_ascii=False).encode('utf-8')
|
||||
data=json.dumps(wechat_message, ensure_ascii=False).encode('utf-8'),
|
||||
timeout=5
|
||||
)
|
||||
|
||||
return True
|
||||
|
||||
Reference in New Issue
Block a user