Sockets, padrões de troca de mensagens e pipelines.
listen nem accept.

AMQP na prática, do Hello World aos exchanges.


$ docker run -d --hostname rabbitmq --name rabbitmq \
-p 5672:5672 -p 15672:15672 \
rabbitmq:3-management
A tag management já vem com o painel web habilitado na porta 15672. hub.docker.com/_/rabbitmq

15672) mostra conexões, canais, exchanges, filas e consumidores em tempo real.
import pika
# Conectando com o rabbit
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# Criando 1 fila
channel.queue_declare(queue='hello')
# Enviando 1 msg para fila
channel.basic_publish(exchange='',
routing_key='hello', # nome da fila
body='Hello World!')
print(" [x] Sent 'Hello World!'")
connection.close()
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
channel.queue_declare(queue='hello')
def callback(ch, method, properties, body):
print(" [x] Received %r" % body)
channel.basic_consume(queue='hello', on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()

routing_key.
routing_key exata — a mensagem só vai para quem casar exatamente com ela.
routing_key aceita padrões com * (uma palavra) e # (zero ou mais palavras).
routing_key — a mensagem vai para todas as filas ligadas ao exchange.Um único evento novo_pedido publicado no fanout — cada fila decide o que fazer com ele.
emit_log.py)
import pika
import sys
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
channel.exchange_declare(exchange='logs', exchange_type='fanout')
message = ' '.join(sys.argv[1:]) or "info: Hello World!"
channel.basic_publish(exchange='logs', routing_key='', body=message)
print(" [x] Sent %r" % message)
connection.close()
Declara um exchange fanout — ele faz o broadcast da mensagem para todas as filas ligadas a ele.
receive_logs.py)
import pika
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
channel.exchange_declare(exchange='logs', exchange_type='fanout')
result = channel.queue_declare(queue='', exclusive=True)
queue_name = result.method.queue
channel.queue_bind(exchange='logs', queue=queue_name)
print(' [*] Waiting for logs. To exit press CTRL+C')
Uma fila com nome aleatório e temporário (exclusive=True) é criada e ligada (binding) ao exchange.
routing_key: alerta — só os alarmes chegam ao Monitor.
Toda leitura chega ao Monitor e ao Climatizador.