|
| 1 | +#!/usr/bin/env python |
| 2 | + |
| 3 | +from pvapy.utility.loggingManager import LoggingManager |
| 4 | +from pvapy.hpc.userMpDataProcessor import UserMpDataProcessor |
| 5 | +from pvapy.hpc.userMpWorkerController import UserMpWorkerController |
| 6 | +import multiprocessing as mp |
| 7 | + |
| 8 | +# Example for implementing data processor that spawns separate unix process |
| 9 | +class UmpDataProcessor2(UserMpDataProcessor): |
| 10 | + def __init__(self): |
| 11 | + UserMpDataProcessor.__init__(self) |
| 12 | + self.udp = UmpDataProcessor() |
| 13 | + self.iq = mp.Queue() |
| 14 | + self.uwpc = UserMpWorkerController(2, self.udp, self.iq) |
| 15 | + |
| 16 | + def start(self): |
| 17 | + self.uwpc.start() |
| 18 | + |
| 19 | + def configure(self, configDict): |
| 20 | + self.configure(configDict) |
| 21 | + |
| 22 | + def process(self, pvObject): |
| 23 | + self.iq.put(pvObject) |
| 24 | + return pvObject |
| 25 | + |
| 26 | + def resetStats(self): |
| 27 | + self.uwpc.resetStats() |
| 28 | + |
| 29 | + def getStats(self): |
| 30 | + return self.uwpc.getStats() |
| 31 | + |
| 32 | + def stop(self): |
| 33 | + self.uwpc.stop() |
| 34 | + |
| 35 | +class UmpDataProcessor(UserMpDataProcessor): |
| 36 | + |
| 37 | + def __init__(self): |
| 38 | + UserMpDataProcessor.__init__(self) |
| 39 | + self.nProcessed = 0 |
| 40 | + |
| 41 | + # Process monitor update |
| 42 | + def process(self, pvObject): |
| 43 | + self.nProcessed += 1 |
| 44 | + self.logger.debug(f'Processing: {pvObject} (nProcessed: {self.nProcessed})') |
| 45 | + return pvObject |
| 46 | + |
| 47 | + # Reset statistics for user processor |
| 48 | + def resetStats(self): |
| 49 | + self.nProcessed = 0 |
| 50 | + |
| 51 | + # Retrieve statistics for user processor |
| 52 | + def getStats(self): |
| 53 | + return {'nProcessed' : self.nProcessed} |
| 54 | + |
| 55 | +if __name__ == '__main__': |
| 56 | + LoggingManager.addStreamHandler() |
| 57 | + LoggingManager.setLogLevel('DEBUG') |
| 58 | + udp = UmpDataProcessor2() |
| 59 | + iq = mp.Queue() |
| 60 | + uwpc = UserMpWorkerController(1, udp, iq) |
| 61 | + uwpc.start() |
| 62 | + import time |
| 63 | + for i in range(0,10): |
| 64 | + iq.put(i) |
| 65 | + time.sleep(1) |
| 66 | + print(uwpc.getStats()) |
| 67 | + statsDict = uwpc.stop() |
| 68 | + print(statsDict) |
0 commit comments