|
|
@ -1,3 +1,4 @@
|
|
|
|
|
|
|
|
################ 待整理
|
|
|
|
'''
|
|
|
|
'''
|
|
|
|
多线程各个模块比较乱的但是协作序贯的完成了数据处理
|
|
|
|
多线程各个模块比较乱的但是协作序贯的完成了数据处理
|
|
|
|
各个组件完全不能互操作,仅依靠队列发消息进行协作
|
|
|
|
各个组件完全不能互操作,仅依靠队列发消息进行协作
|
|
|
@ -10,15 +11,15 @@ from threading import Thread
|
|
|
|
from queue import Queue
|
|
|
|
from queue import Queue
|
|
|
|
from cppy.cp_util import *
|
|
|
|
from cppy.cp_util import *
|
|
|
|
|
|
|
|
|
|
|
|
class ActiveWFObject(Thread):
|
|
|
|
class ThreadObject(Thread):
|
|
|
|
def __init__(self):
|
|
|
|
def __init__(self):
|
|
|
|
super().__init__()
|
|
|
|
super().__init__()
|
|
|
|
self.queue = Queue()
|
|
|
|
self.queue = Queue()
|
|
|
|
self._stopMe = False
|
|
|
|
self._over = False
|
|
|
|
self.start()
|
|
|
|
self.start()
|
|
|
|
|
|
|
|
|
|
|
|
def run(self):
|
|
|
|
def run(self):
|
|
|
|
while not self._stopMe:
|
|
|
|
while not self._over:
|
|
|
|
message = self.queue.get()
|
|
|
|
message = self.queue.get()
|
|
|
|
self._dispatch(message)
|
|
|
|
self._dispatch(message)
|
|
|
|
if message[0] == 'over':
|
|
|
|
if message[0] == 'over':
|
|
|
@ -27,8 +28,7 @@ class ActiveWFObject(Thread):
|
|
|
|
def send(receiver, message):
|
|
|
|
def send(receiver, message):
|
|
|
|
receiver.queue.put(message)
|
|
|
|
receiver.queue.put(message)
|
|
|
|
|
|
|
|
|
|
|
|
class DataStorageManager(ActiveWFObject):
|
|
|
|
class TxtManager(ThreadObject):
|
|
|
|
""" Models the contents of the file """
|
|
|
|
|
|
|
|
_data = ''
|
|
|
|
_data = ''
|
|
|
|
|
|
|
|
|
|
|
|
def _dispatch(self, message):
|
|
|
|
def _dispatch(self, message):
|
|
|
@ -36,8 +36,7 @@ class DataStorageManager(ActiveWFObject):
|
|
|
|
self._init(message[1:])
|
|
|
|
self._init(message[1:])
|
|
|
|
elif message[0] == 'send_word_freqs':
|
|
|
|
elif message[0] == 'send_word_freqs':
|
|
|
|
self._process_words(message[1:])
|
|
|
|
self._process_words(message[1:])
|
|
|
|
else:
|
|
|
|
else:
|
|
|
|
# forward
|
|
|
|
|
|
|
|
send(self._stop_word_manager, message)
|
|
|
|
send(self._stop_word_manager, message)
|
|
|
|
|
|
|
|
|
|
|
|
def _init(self, message):
|
|
|
|
def _init(self, message):
|
|
|
@ -50,8 +49,8 @@ class DataStorageManager(ActiveWFObject):
|
|
|
|
send(self._stop_word_manager, ['filter', w])
|
|
|
|
send(self._stop_word_manager, ['filter', w])
|
|
|
|
send(self._stop_word_manager, ['topWord', recipient])
|
|
|
|
send(self._stop_word_manager, ['topWord', recipient])
|
|
|
|
|
|
|
|
|
|
|
|
class StopWordManager(ActiveWFObject):
|
|
|
|
|
|
|
|
""" Models the stop word filter """
|
|
|
|
class FilterManager(ThreadObject):
|
|
|
|
_stop_words = []
|
|
|
|
_stop_words = []
|
|
|
|
|
|
|
|
|
|
|
|
def _dispatch(self, message):
|
|
|
|
def _dispatch(self, message):
|
|
|
@ -59,8 +58,7 @@ class StopWordManager(ActiveWFObject):
|
|
|
|
self._init(message[1:])
|
|
|
|
self._init(message[1:])
|
|
|
|
elif message[0] == 'filter':
|
|
|
|
elif message[0] == 'filter':
|
|
|
|
return self._filter(message[1:])
|
|
|
|
return self._filter(message[1:])
|
|
|
|
else:
|
|
|
|
else:
|
|
|
|
# forward
|
|
|
|
|
|
|
|
send(self._word_freqs_manager, message)
|
|
|
|
send(self._word_freqs_manager, message)
|
|
|
|
|
|
|
|
|
|
|
|
def _init(self, message):
|
|
|
|
def _init(self, message):
|
|
|
@ -72,8 +70,7 @@ class StopWordManager(ActiveWFObject):
|
|
|
|
if word not in self._stop_words:
|
|
|
|
if word not in self._stop_words:
|
|
|
|
send(self._word_freqs_manager, ['word', word])
|
|
|
|
send(self._word_freqs_manager, ['word', word])
|
|
|
|
|
|
|
|
|
|
|
|
class WordFrequencyManager(ActiveWFObject):
|
|
|
|
class WFManager(ThreadObject):
|
|
|
|
""" Keeps the word frequency data """
|
|
|
|
|
|
|
|
_word_freqs = {}
|
|
|
|
_word_freqs = {}
|
|
|
|
|
|
|
|
|
|
|
|
def _dispatch(self, message):
|
|
|
|
def _dispatch(self, message):
|
|
|
@ -91,8 +88,7 @@ class WordFrequencyManager(ActiveWFObject):
|
|
|
|
freqs_sorted = sort_dict ( self._word_freqs )
|
|
|
|
freqs_sorted = sort_dict ( self._word_freqs )
|
|
|
|
send(recipient, ['topWord', freqs_sorted])
|
|
|
|
send(recipient, ['topWord', freqs_sorted])
|
|
|
|
|
|
|
|
|
|
|
|
class MyController(ActiveWFObject):
|
|
|
|
class MyController(ThreadObject):
|
|
|
|
|
|
|
|
|
|
|
|
def _dispatch(self, message):
|
|
|
|
def _dispatch(self, message):
|
|
|
|
if message[0] == 'run':
|
|
|
|
if message[0] == 'run':
|
|
|
|
self._run(message[1:])
|
|
|
|
self._run(message[1:])
|
|
|
@ -109,13 +105,13 @@ class MyController(ActiveWFObject):
|
|
|
|
word_freqs, = message
|
|
|
|
word_freqs, = message
|
|
|
|
print_word_freqs( word_freqs)
|
|
|
|
print_word_freqs( word_freqs)
|
|
|
|
send(self._storage_manager, ['over'])
|
|
|
|
send(self._storage_manager, ['over'])
|
|
|
|
self._stopMe = True
|
|
|
|
self._over = True
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if __name__ == '__main__':
|
|
|
|
if __name__ == '__main__':
|
|
|
|
word_freq_manager = WordFrequencyManager()
|
|
|
|
word_freq_manager = WFManager()
|
|
|
|
stop_word_manager = StopWordManager()
|
|
|
|
stop_word_manager = FilterManager()
|
|
|
|
storage_manager = DataStorageManager()
|
|
|
|
storage_manager = TxtManager()
|
|
|
|
wfcontroller = MyController()
|
|
|
|
wfcontroller = MyController()
|
|
|
|
|
|
|
|
|
|
|
|
send(storage_manager, ['init', testfilepath, stop_word_manager])
|
|
|
|
send(storage_manager, ['init', testfilepath, stop_word_manager])
|
|
|
|