You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
hbyd_ueba/cron/ueba_cron_pg.py

74 lines
2.6 KiB

3 months ago
# coding=utf-8
"""
3 months ago
@Author: tangwy
3 months ago
@FileName: ueba_cron_pg.py
3 months ago
@DateTime: 2024/7/09 14:19
@Description: 定时清洗es数据
3 months ago
"""
from __future__ import unicode_literals
import random,string
3 months ago
import traceback,json
3 months ago
import time,threading
3 months ago
from uebaMetricsAnalysis.utils.ext_logging import logger_cron
3 months ago
from uebaMetricsAnalysis.utils.db2json import DBUtils, DBType
from uebaMetricsAnalysis.utils.base_dataclean_pg import entry
3 months ago
JOB_STATUS ={
"RUNNING":1,
"FINISH":2,
"ERROR":3
}
3 months ago
class DataCleanCron:
3 months ago
#生成job_id
3 months ago
def generate_job_id(self):
3 months ago
timestamp = int(time.time() * 1000)
random_letters = ''.join(random.choice(string.ascii_letters) for _ in range(7))
return str(timestamp) + random_letters
#每5分钟执行一次
def processing(self):
3 months ago
logger_cron.info("JOB:接收到执行指令")
3 months ago
job_id =self.generate_job_id()
3 months ago
task_run_count =0
3 months ago
try:
3 months ago
start,end,status,run_count,jobid= DBUtils.get_job_period()
if jobid !="":
job_id=jobid
3 months ago
if end<start:
logger_cron.info("JOB:"+job_id+"开始时间大于结束时间不执行")
return
3 months ago
logger_cron.info("JOB:"+job_id+"开始执行")
3 months ago
if status ==1:
logger_cron.info("JOB:"+job_id+"正在运行中不执行")
return
3 months ago
#延迟15分钟读取es数据
3 months ago
if start is None or end is None:
3 months ago
logger_cron.info("JOB:"+job_id+"结束时间大于(服务器时间-15分钟)不执行")
3 months ago
return
3 months ago
task_run_count = run_count+1
3 months ago
logger_cron.info("JOB:"+job_id+"运行参数:{},{}".format(start,end))
logger_cron.info("JOB:"+job_id+"准备将job写入job表")
3 months ago
DBUtils.insert_job_record(job_id,start,end,JOB_STATUS.get("RUNNING"))
3 months ago
logger_cron.info("JOB:"+job_id+"完成job表写入")
3 months ago
3 months ago
logger_cron.info("JOB:"+job_id+"准备获取es数据")
entry(start,end,job_id)
logger_cron.info("JOB:"+job_id+"完成es数据获取")
3 months ago
DBUtils.write_job_status(job_id,JOB_STATUS.get("FINISH"),"",task_run_count)
3 months ago
logger_cron.info("JOB:"+job_id+"更新job表状态完成")
3 months ago
3 months ago
except Exception ,e:
3 months ago
err_info=traceback.format_exc()
logger_cron.error("JOB:"+job_id+"执行失败:"+err_info)
3 months ago
DBUtils.write_job_status(job_id,JOB_STATUS.get("ERROR"),err_info,task_run_count)
3 months ago
raise
if __name__ == '__main__':
3 months ago
DataCleanCron().processing()