I'm using ZeroMQ to establish a publisher/subscriber communication model.
The publisher creates a zmq context and then opens a socket with the PUB
communication pattern. It then binds to a port, as the TCP transport protocol is used. For the synchronization a separate socket is opened with the REP communication pattern that binds in a different path. Unless a synchronization request in msg = syncservice.recv() is received, the program cannot continue. It then does some rudimentary work and starts over again. Here's the code for the publisher:
import pickle, zmq, random, string
# Wait for 1 subscriber
SUBSCRIBERS_EXPECTED = 1
def randomword(length):
letters = string.ascii_lowercase
return ''.join(random.choice(letters) for i in range(length))
while True:
try:
arguments = {}
data = {}
context = zmq.Context()
# Socket to talk to clients
publisher = context.socket(zmq.PUB)
# set SNDHWM, in case of slow subscribers
publisher.sndhwm = 1100000
publisher.bind('tcp://*:5561')
# Socket to receive signals
syncservice = context.socket(zmq.REP)
syncservice.bind('tcp://*:5562')
# Get synchronization from subscribers
subscribers = 0
while subscribers < SUBSCRIBERS_EXPECTED:
# wait for synchronization request
msg = syncservice.recv()
# send synchronization reply
syncservice.send(b'')
subscribers += 1
for n in range(1000):
for i in range(random.randrange(1, 6)):
arguments[i] = randomword(random.randrange(2, 10))
data['func_name_' + str(n)] = randomword(8)
data['arguments_' + str(n)] = arguments
data_string = pickle.dumps(data)
publisher.send(data_string)
except KeyboardInterrupt:
print("Interrupt received, stopping...")
break
The subscriber functions pretty much the same way the publisher does, albeit from a subscriber perspective. Here's the code for the subscriber:
import pickle, zmq, pprint, time
context = zmq.Context()
# Connect the subscriber socket
subscriber = context.socket(zmq.SUB)
subscriber.connect('tcp://localhost:5561')
subscriber.setsockopt(zmq.SUBSCRIBE, b'')
time.sleep(1)
# Synchronize with publisher
syncclient = context.socket(zmq.REQ)
syncclient.connect('tcp://localhost:5562')
# Initialize poll set
poller = zmq.Poller()
poller.register(syncclient, zmq.POLLIN)
poller.register(subscriber, zmq.POLLIN)
# send a synchronization request
syncclient.send(b'')
while True:
try:
socks = dict(poller.poll())
except KeyboardInterrupt:
print("Interrupt received, stopping...")
break
# wait for synchronization reply
if syncclient in socks:
syncclient.recv()
print('Sync')
if subscriber in socks:
msg = subscriber.recv()
data = pickle.loads(msg)
pprint.pprint(data)
syncclient.send(b'')
The desired result would be for the publisher to publish endlessly, while the subscriber continually receives and prints everything. If I remove the synchronization part, everything runs as expected. If I keep the synchronization part the subscriber hangs after a number of transmissions. The interesting thing is that if I send a keyboard interrupt (Ctrl-C) and then restart the subscriber, it will receive a couple of transmissions again and hang again and so on and so forth.
I have tried different high-watermark settings, but it didn't make any difference. I have tried closing the sockets and terminating the context after every loop. I've tested if the overhead from printing or pickling (serializing) was too much, but it wasn't it either. I have also modified the suicidal snail example to work in this case, but the subscriber didn't die. What am I missing? (Python 3 is used for every example)