· 8 years ago · May 14, 2018, 10:30 PM
1# -*- coding:utf8 -*-
2"""
3Author : Myth
4Date : 2018/5/9
5Email : email4myth at gmail.com
6"""
7
8import json
9import os
10import signal
11import time
12from collections import OrderedDict
13from multiprocessing import Process, Queue
14
15from confluent_kafka import Consumer
16from datetime import datetime
17
18bro_group_id = 'myth.worker.count.test'
19queue = Queue(100)
20
21consumers = OrderedDict()
22
23'''
24{
25 pid: {
26 'create_time': TIME,
27 'last_update': TIME
28 }
29}
30'''
31
32
33def consumer(queue):
34 c = Consumer({'group.id': bro_group_id, 'enable.auto.commit': True, 'bootstrap.servers': '10.0.81.9:9091'})
35 c.subscribe(['test-topic'])
36 while True:
37 msg = c.poll(3)
38 if not msg:
39 continue
40 if not msg.error():
41 v = msg.value()
42 queue.put(json.dumps({'pid': str(os.getpid()), 'ts': time.time()}))
43
44
45def handle_signal(s, f):
46 if s == signal.SIGHUP:
47 start_new_consumer()
48 else:
49 print 'stoping all sub processes ...'
50 for pid in consumers:
51 pid_info = consumers[pid]
52 p = pid_info['p']
53 p.terminate()
54
55
56def check_consumers():
57 while True:
58
59 # update lost update
60 while not queue.empty():
61 last_info = json.loads(queue.get())
62 pid = last_info['pid']
63 ts = last_info['ts']
64 consumers[pid]['last_update'] = datetime.fromtimestamp(ts)
65
66 os.system('clear')
67 pids = consumers.keys()
68 print ' pid | ctime | utime | duration (%s consumers)' % len(pids)
69 print '-----------------------------------------------------------'
70 for pid in pids:
71 consumer_info = consumers[pid]
72 now = datetime.now()
73 create_time = consumer_info['create_time']
74 last_update = consumer_info['last_update']
75 if not last_update:
76 last_update_str = '<notime>'
77 duration = 0
78 else:
79 duration = (now - last_update).total_seconds()
80 last_update_str = last_update.strftime('%H:%M:%S')
81
82 print '%s | %s | %s | %s %s' % (pid.zfill(5), create_time.strftime('%H:%M:%S'), last_update_str, duration,
83 '<<<<<<<<<<<<<<<<<<<' if duration > 30 else '')
84
85 time.sleep(3)
86
87
88def start_new_consumer():
89 p = Process(target=consumer, args=(queue,))
90 p.start()
91 consumers[str(p.pid)] = {
92 'p': p,
93 'create_time': datetime.now(),
94 'last_update': None
95 }
96
97
98def main():
99 for i in range(30):
100 start_new_consumer()
101
102 check_consumers()
103
104
105if __name__ == '__main__':
106 signal.signal(signal.SIGINT, handle_signal)
107 signal.signal(signal.SIGTERM, handle_signal)
108 signal.signal(signal.SIGHUP, handle_signal)
109 main()