PyKafka has the limitation that:
delivery report queue is thread-local: it will only serve reports for messages which were produced from the current thread
I'm trying to write a script where I can asynchronously send messages using one function, and keep receiving acknowledgments via another function.
Here are the functions:
def SendRequest(producer):
count=0
while True:
count += 1
producer.produce('test msg', partition_key='{}'.format(count))
if count == 50000:
endtime=datetime.datetime.now()
print "EndTime : ",endtime
print "Done sending all messages.Waiting for response now"
return
def GetResponse(producer):
count_response=0
while True:
try:
msg, exc = producer.get_delivery_report(block=False)
if exc is not None:
count_response+=1
print 'Failed to deliver msg {}: {}'.format(
msg.partition_key, repr(exc))
else:
print "Count Res :",count_response
count_response+=1
except Queue.Empty:
pass
except Exception,e:
print "Unhandled exception : ",e
Threading and multiprocessing did not help. These above two functions need to be running asynchronously/in parallel. What approach shall be used here?