-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathThreadMgr.py
More file actions
45 lines (38 loc) · 1.24 KB
/
Copy pathThreadMgr.py
File metadata and controls
45 lines (38 loc) · 1.24 KB
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
#!/usr/bin/python
#encoding=utf-8
import os
import sys
import time
import datetime
import logging
import redis
from Thread import GPSThread
import threading
class ThreadMgr:
def __init__(self, thread_num=1, logger=None):
self.thread_num = thread_num
self.threads = []
self.queues = []
self.mutexs = []
self.logger = logger
for i in range(self.thread_num):
self.queues.append( [] )
self.mutexs.append( threading.Lock() )
self.threads.append( GPSThread(i, self.queues[i], self.mutexs[i], logger) )
def push_record(self, record):
thread_id = self.hash_id(record.host_fd) % self.thread_num
self.mutexs[thread_id].acquire()
self.queues[thread_id].append(record)
self.mutexs[thread_id].release()
print 'thread[{0}] process the record[{1}]'.format(thread_id, record.data)
return True
def process(self):
for i in range(self.thread_num):
self.threads[i].start()
def end_process(self):
for i in range(self.thread_num):
if self.threads[i] is threading.currentThread():
continue
self.threads[i].join()
def hash_id(self, g_id):
return int(g_id)