python工程,kafka消息对接转换成配置文件并完成类似看门狗监控linux系统进程,定期清理日志
wangrong
2025-04-05 88f9100edcba9581214c1a0acc23cd33a821a9c6
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
import configparser
import json
import psutil
import subprocess
import time
import threading
import logging
import os
 
 
# 设置日志记录
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)
 
def read_msg_config():
    """读取 msg_config.ini 配置文件"""
    config = configparser.ConfigParser()
    
    if not os.path.exists('msg_config.ini'):
        logger.error("msg_config.ini 文件不存在")
        raise FileNotFoundError("msg_config.ini 文件不存在")
    
    config.read('msg_config.ini')
    
    try:
        kafka_server = config.get('KAFKA', 'server')
        kafka_topics = config.get('KAFKA', 'topics').split(', ')
        process_file = config.get('PROCESSES', 'process_file')
    except configparser.NoOptionError as e:
        logger.error(f"配置文件格式错误: {e}")
        raise
    
    return kafka_server, kafka_topics, process_file
 
 
def read_process_list(process_file):
    """读取 process_list.json 配置文件"""
    if not os.path.exists(process_file):
        logger.error(f"{process_file} 文件不存在")
        raise FileNotFoundError(f"{process_file} 文件不存在")
    
    with open(process_file, 'r') as f:
        try:
            process_data = json.load(f)
        except json.JSONDecodeError as e:
            logger.error(f"JSON 格式错误: {e}")
            raise
    return process_data
 
 
def stop_processes(prefixes):
    """停止符合进程名以任意prefix打头的进程"""
    for proc in psutil.process_iter(attrs=['pid', 'name']):
        try:
            # 检查进程名是否以任何一个prefix打头
            if any(proc.info['name'].lower().startswith(prefix.lower()) for prefix in prefixes):
                proc.kill()
                logger.info(f"Terminated process: {proc.info['name']} (PID: {proc.info['pid']})")
        except (psutil.NoSuchProcess, psutil.AccessDenied):
            continue    
 
 
def start_process(proc):
    """启动单个进程"""
    program_path = proc.get('program_path')
    command_ = proc.get('commands')
    
    try:
        if program_path == "gnome-terminal":
            subprocess.Popen([program_path, '--', 'bash', '-c', command_])
        elif program_path == "xterm":
            subprocess.Popen([program_path, '-e', 'bash -c "' + command_ + '"'])
        elif program_path == "konsole":
            subprocess.Popen([program_path, '--hold', '-e', 'bash -c "' + command_ + '"'])
        else:
            subprocess.Popen([program_path, '-c', command_])
        
        logger.info(f"Started process: {command_}")
    except FileNotFoundError as e:
        logger.error(f"终端程序未找到: {e}")
    except Exception as e:
        logger.error(f"启动进程失败: {e}")
 
 
def start_processes_concurrently(process_list):
    """并行启动多个进程"""
    threads = []
    
    for proc in process_list:
        thread = threading.Thread(target=start_process, args=(proc,))
        threads.append(thread)
        thread.start()
    
    for thread in threads:
        thread.join()  # 等待所有线程完成
 
 
if __name__ == "__main__":
    # 读取配置文件
    try:
        kafka_server, kafka_topics, process_file = read_msg_config()
        logger.info(f"KAFKA Server: {kafka_server}")
        logger.info(f"KAFKA Topics: {kafka_topics}")
    except Exception as e:
        logger.error(f"读取配置文件失败: {e}")
        exit(1)
    
    # 读取进程清单
    try:
        process_data = read_process_list(process_file)
        process_names = process_data['process_names']
        process_list = process_data['process_list']
    except Exception as e:
        logger.error(f"读取进程列表失败: {e}")
        exit(1)
    
    # 停止符合进程名的所有进程
    stop_processes(process_names)
    
    # 启动新进程
    start_processes_concurrently(process_list)
 
 
'''
用python写一份程序,实现以下功能:
1. 读取msg_config.ini文件,msg_config.ini文件内容如下:
[KAFKA]
server = 192.168.0.79:9092
topics = vp_ba_jam, vp_ba_stop, vp_ba_crossline, vp_tracking
 
[PROCESSES]
process_file = textjson/process_list.json
 
2. 获取process_file指定的文件路径,读取process_list.json文件内容,如下:
{
    "process_names": ["top", "mixc2", "mixc3"],
    "process_list": [
        {"program_path": "gnome-terminal", "commands": "top -d 1"},
        {"program_path": "gnome-terminal", "commands": "top -p 1234"},
        {"program_path": "gnome-terminal", "commands": "top -u wr"},   
        {"program_path": "gnome-terminal", "commands": "top -H"} 
    ]
}
3. 读取process_names,停止符合进程名的所有进程;
 
4. 停止上述进程之后,读取process_list进程清单,并调起命令行终端启动process_list中的所有程序,每隔5秒启动一个,直到所有程序启动完毕;
 
'''