Commit 25b2a077 authored by payno's avatar payno
Browse files

[JoinActor] adapt join actor to updated code

parent b88d79cf
......@@ -36,14 +36,18 @@ class JoinActor(AbstractActor):
def increaseNumberOfThreads(self):
self.numberOfThreads += 1
def trigger(self, inData):
self.listInData.append(inData)
def trigger(self, in_data):
if in_data is None:
channel = data = None
else:
channel, data = in_data
self.listInData.append(data)
if len(self.listInData) == self.numberOfThreads:
newInData = {}
for data in self.listInData:
newInData.update(data)
for actor in self.listDownStreamActor:
actor.trigger(newInData)
actor.trigger((channel, data))
class JoinUntilStopSignal(AbstractActor):
......
Supports Markdown
0% or .
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment