Programação Concorrente e Distribuída

Mensageria e AMQP

Aula 12 — Padrões de mensagens e RabbitMQ na prática
Prof. Newton S. B. Miyoshi
newton.miyoshi@baraodemaua.br
Centro Universitário Barão de Mauá — Ribeirão Preto
Parte 1

Comunicação Orientada a Mensagens

Sockets, padrões de troca de mensagens e pipelines.

Padrões de Mensagens

  • A comunicação ocorre pelo pareamento de sockets: um tipo específico de socket para envio é pareado com o tipo correspondente para recebimento.
  • Cada par de tipos de socket corresponde a um padrão de comunicação.
  • Três padrões mais comuns:
    • request-reply
    • publish-subscribe
    • pipeline

Padrão Request-Reply

  • Usado em comunicações tradicionais cliente-servidor / RPC.
  • Cliente utiliza um request socket (REQ).
  • Servidor utiliza um reply socket (REP).
  • Não é necessário chamar as operações listen nem accept.
ClienteREQ
ServidorREP

Padrão Publish-Subscribe

  • Clientes subscrevem mensagens específicas que são publicadas no servidor.
  • Somente as mensagens subscritas são notificadas.
  • Base para sistemas orientados a eventos.
  • Implementa um modelo multicast de comunicação.
Publisher envia mensagem para o EventCenter, que notifica os Subscribers que se inscreveram
Subscribers se inscrevem (subscribe) no evento; quando o publisher envia uma mensagem, todos os inscritos são notificados (notify).

Padrão Pipeline

  • Nós/processos empurram (push) mensagens para um duto.
  • Outros nós/processos puxam (pull) mensagens do duto.
  • Quem empurra não se importa quem vai puxar a mensagem — e vice-versa.
  • Objetivo: performance no fluxo de mensagens.
Vários Producers empurram mensagens para uma fila; vários Consumers puxam mensagens da mesma fila
Vários producers empurram mensagens para uma fila; vários consumers puxam dessa mesma fila — sem se conhecerem.
Parte 2

Python com RabbitMQ

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

AMQP — visão geral

Publisher publica mensagens num Exchange, que roteia para Queues com base em bindings; Consumers assinam ou pedem mensagens das queues
O publisher nunca envia direto para uma fila — ele publica num exchange, que roteia a mensagem para as queues ligadas a ele.
Para discutir
?

No projeto Café com Pão, quais seriam possíveis eventos gerados?

Notações — RabbitMQ

PProducer
queue_nameQueue
CConsumer

RabbitMQ — Docker

Quick reference da imagem oficial do RabbitMQ no Docker Hub, com as tags suportadas

$ 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

RabbitMQ Management

Painel de administração web do RabbitMQ, aba Overview, mostrando totais de conexões, canais, exchanges, filas e consumidores
O painel web (porta 15672) mostra conexões, canais, exchanges, filas e consumidores em tempo real.

Hello World — enviando


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()
        

Hello World — recebendo


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()
        

Pub-Sub com RabbitMQ

Producer envia para exchange X, que roteia para duas filas com nomes gerados automaticamente, cada uma lida por um consumer diferente
Um exchange X pode rotear a mesma publicação para várias filas — cada uma com seu próprio consumer.

Exchanges

  • Producer é a aplicação que envia mensagens; Queue é um buffer que as armazena; Consumer é quem recebe.
  • O producer não envia direto para a fila — geralmente ele nem sabe se/quando a mensagem foi entregue.
  • O producer envia para uma Exchange, que pode selecionar a fila via routing_key.
  • Tipos de exchange: direct e fanout (outros: topic e header).

Exchange — Direct

Exchange direct roteando order-create para order_create_queue e order-create-log para order_create_log_queue
Cada fila se liga a uma routing_key exata — a mensagem só vai para quem casar exatamente com ela.

Exchange — Topic

Exchange topic roteando por padrões como order.logs.customer.#, order.logs.# e order.logs.*.electronics
A routing_key aceita padrões com * (uma palavra) e # (zero ou mais palavras).

Exchange — Fanout

Exchange fanout enviando a mesma mensagem para warehouse_queue, cargo_queue e logs_queue, ignorando routing key
Ignora a routing_key — a mensagem vai para todas as filas ligadas ao exchange.
Café com Pão na prática

Um novo pedido, vários interessados

Novo pedido
Café com Pãofanout
Mailing
WhatsApp
Notificação
Entregas
NF-e
Sistema de Notas
Entregadores×5, cada um consome da mesma fila

Um único evento novo_pedido publicado no fanout — cada fila decide o que fazer com ele.

Fanout — publicando logs (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.

Fanout — assinando logs (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.

Direct × Fanout — uma rede de sensores

Sensor
Sensor
Sensor
↓ filtro

Direct

routing_key: alerta — só os alarmes chegam ao Monitor.

Fanout

Toda leitura chega ao Monitor e ao Climatizador.

Referências bibliográficas