I have a RabbitMQ server and a Flask server that consumes a RabbitMQ queue. Messages are published by another server. Clients connect to the Flask server using SocketIO, and on connection I store the session ID in a synchronized queue. The message consumption is done in a a separate thread. When a message arrives, the callback function takes the first item in the queue and sends a SocketIO message to the respective client.
I tried to implement all this, but when a message arrives and the callback function is fired, the queue is empty, while before the message arrived, I connected a client to it and verified that indeed its session id is put into the queue.
My guess is that the consumer thread uses a copy of the queue which is empty when the thread starts, and after the queue is updated in the Flask thread after a client connects, this change is not applied in the consumer thread.
This is my code:
app = Flask(__name__)
socketio = SocketIO(app, ping_interval=5, async_mode='threading')
socketio.init_app(app, cors_allowed_origins="*")
client_queue = queue.Queue()
@socketio.on('connect')
def client_connected():
client_queue.put(request.sid, block=False)
def callback(ch, method, properties, body):
try:
selected_client = client_queue.get(block=True, timeout=5)
except queue.Empty as e:
print(e)
print("No clients")
connection = pika.BlockingConnection(pika.ConnectionParameters(host=os.getenv("RABBITMQ")))
channel = connection.channel()
channel.queue_declare(queue=os.getenv("QUEUE_NAME"), durable=True)
channel.basic_consume(queue=os.getenv("QUEUE_NAME"), on_message_callback=callback)
if __name__ == '__main__':
socketio.start_background_task(target=channel.start_consuming)
socketio.run(app, debug=True)
So when I run this code, connect to the server with a SocketIO client, and after that send a message through RabbitMQ, it will result in the Empty queue expection (printing no clients). What am I doing wrong here?