feat:增加批量插入和查询日志告警功能
This commit is contained in:
BIN
service/__pycache__/log_alert_service.cpython-310.pyc
Normal file
BIN
service/__pycache__/log_alert_service.cpython-310.pyc
Normal file
Binary file not shown.
BIN
service/__pycache__/message_service.cpython-310.pyc
Normal file
BIN
service/__pycache__/message_service.cpython-310.pyc
Normal file
Binary file not shown.
BIN
service/__pycache__/openobserve_service.cpython-310.pyc
Normal file
BIN
service/__pycache__/openobserve_service.cpython-310.pyc
Normal file
Binary file not shown.
@@ -1,31 +1,70 @@
|
||||
import requests
|
||||
from requests.auth import HTTPBasicAuth
|
||||
from collections import defaultdict
|
||||
|
||||
from schemas.log_alert import LogAlertTrigger, OpenobserveQuery
|
||||
from sqlalchemy import insert, select
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from models.log_alert import LogAlert
|
||||
from schemas.log_alert import LogAlertTrigger
|
||||
from service import openobserve_service
|
||||
from service.message_service import send_wechat_message, MessageEnum
|
||||
from utils.time import get_timestamp_range, format_timestamp
|
||||
|
||||
|
||||
def trigger_alert(log_alter: LogAlertTrigger):
|
||||
#print(log_alter)
|
||||
# openobserve_service.get_log()
|
||||
def trigger_alert(log_alter: LogAlertTrigger, db: Session) -> bool:
|
||||
start_time, end_time = get_timestamp_range(log_alter.timestamp, 60)
|
||||
result = openobserve_service.get_log(start_time, end_time)
|
||||
hits = result.get("hits")
|
||||
|
||||
url = 'http://192.168.1.7:5080/api/default/_search'
|
||||
username = 'njcxx0822@163.com'
|
||||
password = '19940822Cxx'
|
||||
# 判断日志是否为空
|
||||
if len(hits) == 0:
|
||||
return False
|
||||
|
||||
response = requests.post(
|
||||
url,
|
||||
headers={
|
||||
'accept': 'application/json',
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
json={
|
||||
"query": OpenobserveQuery(
|
||||
start_time=1760773184953420,
|
||||
end_time=1760773189953431
|
||||
).model_dump()
|
||||
},
|
||||
auth=HTTPBasicAuth(username, password)
|
||||
# 判断日志告警是否已经存在
|
||||
log_db_alter = db.execute(
|
||||
select(LogAlert)
|
||||
.where(LogAlert.alter_timestamp == log_alter.timestamp)
|
||||
).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)
|
||||
|
||||
|
||||
def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool:
|
||||
stmt = insert(LogAlert).values(
|
||||
[
|
||||
{
|
||||
"alter_name": log_alter.alter_name,
|
||||
"alter_timestamp": log_alter.timestamp,
|
||||
"timestamp": hit["_timestamp"],
|
||||
"message": hit["message"],
|
||||
"status": "pending"
|
||||
}
|
||||
for hit in hits
|
||||
]
|
||||
)
|
||||
|
||||
return response.json().get("hits")
|
||||
db.execute(stmt)
|
||||
db.commit()
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def query_alert(db: Session):
|
||||
alerts = db.execute(select(LogAlert)).scalars().all()
|
||||
|
||||
# 按照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
|
||||
})
|
||||
|
||||
return dict(classified)
|
||||
|
||||
32
service/message_service.py
Normal file
32
service/message_service.py
Normal file
@@ -0,0 +1,32 @@
|
||||
import json
|
||||
from enum import Enum
|
||||
|
||||
import requests
|
||||
|
||||
from config.setting import settings
|
||||
|
||||
|
||||
class MessageEnum(Enum):
|
||||
TEXT = 'text'
|
||||
MARKDOWN = 'markdown'
|
||||
MARKDOWN2 = 'markdown2'
|
||||
|
||||
|
||||
def send_wechat_message(message_type: MessageEnum, message: str):
|
||||
wechat_message = {}
|
||||
match message_type:
|
||||
case MessageEnum.TEXT:
|
||||
wechat_message = {"msgtype": "text", "text": {"content": message}}
|
||||
case MessageEnum.MARKDOWN:
|
||||
wechat_message = {"msgtype": "markdown", "markdown": {"content": message}}
|
||||
case MessageEnum.MARKDOWN2:
|
||||
wechat_message = {"msgtype": "markdown_v2", "markdown_v2": {"content": message}}
|
||||
|
||||
requests.post(
|
||||
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')
|
||||
)
|
||||
|
||||
return True
|
||||
@@ -1,27 +1,19 @@
|
||||
import requests
|
||||
from requests.auth import HTTPBasicAuth
|
||||
|
||||
from config.setting import settings
|
||||
from schemas.log_alert import OpenobserveQuery
|
||||
|
||||
|
||||
def get_log():
|
||||
url = 'http://192.168.1.7:5080/api/default/_search'
|
||||
username = 'njcxx0822@163.com'
|
||||
password = '19940822Cxx'
|
||||
|
||||
def get_log(start_time: int, end_time: int) -> str:
|
||||
response = requests.post(
|
||||
url,
|
||||
settings.openobserve_url,
|
||||
headers={
|
||||
'accept': 'application/json',
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
json={
|
||||
"query": OpenobserveQuery(
|
||||
start_time=1760773184953420,
|
||||
end_time=1760773189953431
|
||||
).model_dump()
|
||||
},
|
||||
auth=HTTPBasicAuth(username, password)
|
||||
json={"query": OpenobserveQuery(start_time=start_time, end_time=end_time).model_dump()},
|
||||
auth=HTTPBasicAuth(settings.openobserve_username, settings.openobserve_password)
|
||||
)
|
||||
|
||||
print(response)
|
||||
return response.json()
|
||||
|
||||
Reference in New Issue
Block a user