Estoy tratando de crear una aplicación basada en RabbitMQ-FastAPI muy simple que pueda ejecutarse en Kubernetes. Se crearon dos API separadas para productor y consumidor. Productor API:
@app.get("/hello/{message}") def post_message(message: str = Body(..., example="Hello World!")): try: connection = pika.BlockingConnection( pika.ConnectionParameters(host="rabbitmq-0.rabbitmq-headless.keda.svc.cluster.local", port=5672, credentials=pika.PlainCredentials("user", "PASSWORD"))) channel = connection.channel() #channel.exchange_declare(exchange='logs', exchange_type='direct') #severity = ['hello', 'message', 'error'] #messages = ['Hafizur', 'message', 'error'] channel.queue_declare(queue='hello') channel.basic_publish(exchange='', routing_key='hello', body=message) connection.close() return {"Successfully sended {} do queue".format(message)}API del consumidor:
@app.post("/message") def get_message(): try: connection = pika.BlockingConnection( pika.ConnectionParameters(host="rabbitmq-0.rabbitmq-headless.keda.svc.cluster.local", port=5672, credentials=pika.PlainCredentials("user", "PASSWORD"))) channel = connection.channel() channel.queue_declare(queue='hello') channel.basic_consume(on_message_callback=callback, queue="hello", auto_ack=True) print(" [*] Waiting for messages. To exit press CTRL+C") channel.start_consuming() except Exception as e: raise HTTPException(status_code=500, detail=str(e))A partir de aquí estoy construyendo dos imágenes acoplables, una para el productor y otra para el consumidor. Usando esos construidos dos implementaciones en K8S. Hizo el reenvío de puertos y accedió a la API del productor e intentó enviar algún mensaje a la cola de RabbitMQ. Pero obteniendo Error: Entidad no procesable.
Estoy buscando su ayuda para resolver este problema.