Guia Prático: Microserviços com DDD, CQRS, Redis e RabbitMQ
Stack: Node.js + TypeScript + Express + Prisma + PostgreSQL + Redis + RabbitMQ + Docker
Nível: Estudo/aprendizagem prática
Objetivo: Construir um mini sistema de Pedidos e Inventário dividido em microserviços, aplicando DDD (Domain-Driven Design), CQRS (Command Query Responsibility Segregation) e comunicação assíncrona via mensageria.
0. Antes de começar: os conceitos que vamos aplicar
Vale a pena entender o quê e o porquê antes do como, senão o código vira "magia".
DDD (Domain-Driven Design)
É uma forma de organizar o código em torno do domínio do negócio, não da tecnologia. Em vez de pastas controllers/, models/, services/ genéricas, organizamos por conceitos do negócio (Pedido, Cliente, Estoque). Ideias-chave que vamos usar:
-
Entidade: objeto com identidade própria que persiste no tempo (ex:
Order, com umid). -
Value Object: objeto sem identidade, definido pelos seus valores (ex:
Money,Address). Dois VOs com os mesmos valores são "iguais". -
Aggregate (Agregado): um conjunto de entidades/VOs tratado como uma unidade de consistência. Tem uma raiz (Aggregate Root) que controla o acesso — ex:
Orderé a raiz deOrderItem. -
Domain Event: algo que aconteceu no domínio e que interessa a outras partes do sistema (ex:
OrderCreated,OrderPaid). - Repository: abstração para persistir/recuperar agregados, escondendo os detalhes do banco de dados.
- Bounded Context: fronteira onde um modelo de domínio é válido. Cada microserviço deste guia é o dono de um Bounded Context (Pedidos, Estoque).
CQRS (Command Query Responsibility Segregation)
Separa o lado de escrita (Commands: "criar pedido", "cancelar pedido") do lado de leitura (Queries: "listar pedidos", "ver detalhes"). Vantagens práticas:
- O modelo de escrita fica focado nas regras de negócio e consistência.
- O modelo de leitura pode ser desnormalizado, otimizado para consulta rápida (e cacheado em Redis), sem se preocupar com regras de negócio.
- Os dois lados podem até viver em bancos de dados diferentes, sincronizados por eventos.
Arquitetura orientada a eventos (Event-Driven) com RabbitMQ
Quando o serviço de Pedidos cria um pedido, ele não chama diretamente o serviço de Estoque (isso seria acoplamento forte, síncrono, frágil). Em vez disso, ele publica um evento (OrderCreated) numa fila do RabbitMQ. O serviço de Estoque assina essa fila e reage de forma assíncrona. Isso dá:
- Desacoplamento entre serviços.
- Resiliência (se o Estoque cair, a mensagem fica na fila até ele voltar).
- Possibilidade de vários serviços reagirem ao mesmo evento.
Redis
Vamos usá-lo de duas formas neste projeto:
- Cache do lado de leitura (CQRS) — respostas de consultas frequentes ficam em cache, reduzindo carga no banco.
- (Opcional, mencionado no final) Idempotência de mensagens — evitar processar o mesmo evento duas vezes.
Consistência eventual
Como os serviços se comunicam por eventos assíncronos, não existe uma transação única que atualiza tudo instantaneamente. Existe um pequeno intervalo em que os dados estão "temporariamente desincronizados" até o evento ser processado. Isso é normal e esperado em sistemas distribuídos — chama-se consistência eventual, e é uma troca consciente por escalabilidade e desacoplamento.
1. Visão geral da arquitetura
Vamos construir 3 microserviços independentes (cada um com seu próprio package.json, seu próprio banco de dados e seu próprio processo):
┌─────────────────────┐
│ Order Service │ (escrita - Commands)
│ Express + Prisma │
│ PostgreSQL (orders) │
└──────────┬───────────┘
│ publica evento
▼
┌─────────────────────┐
│ RabbitMQ │
│ exchange: orders_ex │
└───┬───────────────┬───┘
│ │
consome │ │ consome
▼ ▼
┌─────────────────────┐ ┌──────────────────────┐
│ Inventory Service │ │ Query Service (CQRS) │ (leitura - Queries)
│ Express + Prisma │ │ Express + Redis │
│ PostgreSQL (stock) │ │ PostgreSQL (read-db) │
└─────────────────────┘ └──────────────────────┘
Fluxo de ponta a ponta:
- Cliente faz
POST /ordersno Order Service → cria o pedido no banco de escrita → publica o eventoOrderCreatedno RabbitMQ. - O Inventory Service consome
OrderCreated→ reserva/abate o estoque. - O Query Service também consome
OrderCreated→ atualiza sua própria tabela de leitura (desnormalizada) e invalida/atualiza o cache Redis. - Cliente faz
GET /orders/:idno Query Service → responde rápido, servindo do Redis quando possível.
Cada serviço é dono exclusivo do seu banco de dados — nenhum serviço acessa a tabela de outro diretamente. Essa é a regra de ouro dos microserviços.
2. Pré-requisitos
- Node.js 20+ e npm
- Docker e Docker Compose
- Conhecimento básico de TypeScript, Express e SQL
- Um cliente HTTP para testar (Insomnia, Postman ou
curl)
Instale globalmente (opcional, mas ajuda a inspecionar filas):
npm install -g pnpm
(Vamos usar npm nos exemplos, mas sinta-se livre por usar pnpm.)
3. Estrutura de pastas do projeto (monorepo simples)
Para um projeto de estudo, um monorepo com pastas separadas por serviço é suficiente (sem ferramentas como Nx/Turborepo, para manter o foco no conteúdo):
mini-ecommerce-ms/
├── docker-compose.yml
├── order-service/
│ ├── prisma/
│ │ └── schema.prisma
│ ├── src/
│ │ ├── domain/
│ │ │ ├── entities/
│ │ │ │ └── Order.ts
│ │ │ ├── value-objects/
│ │ │ │ └── Money.ts
│ │ │ ├── events/
│ │ │ │ └── OrderCreated.ts
│ │ │ └── repositories/
│ │ │ └── IOrderRepository.ts
│ │ ├── application/
│ │ │ └── commands/
│ │ │ ├── CreateOrder.command.ts
│ │ │ └── CreateOrder.handler.ts
│ │ ├── infrastructure/
│ │ │ ├── database/
│ │ │ │ └── PrismaOrderRepository.ts
│ │ │ └── messaging/
│ │ │ └── RabbitMQPublisher.ts
│ │ ├── interfaces/
│ │ │ └── http/
│ │ │ ├── order.routes.ts
│ │ │ └── order.controller.ts
│ │ └── server.ts
│ ├── package.json
│ └── tsconfig.json
├── inventory-service/
│ ├── prisma/schema.prisma
│ ├── src/
│ │ ├── domain/...
│ │ ├── application/...
│ │ ├── infrastructure/
│ │ │ └── messaging/
│ │ │ └── RabbitMQConsumer.ts
│ │ └── server.ts
│ ├── package.json
│ └── tsconfig.json
└── query-service/
├── prisma/schema.prisma
├── src/
│ ├── application/
│ │ └── queries/
│ │ ├── GetOrderById.query.ts
│ │ └── GetOrderById.handler.ts
│ ├── infrastructure/
│ │ ├── cache/
│ │ │ └── RedisClient.ts
│ │ └── messaging/
│ │ └── RabbitMQConsumer.ts
│ ├── interfaces/http/
│ └── server.ts
├── package.json
└── tsconfig.json
A ideia: cada serviço replica a mesma organização em camadas (domain → application → infrastructure → interfaces), mesmo tendo responsabilidades diferentes. Isso é a "Arquitetura em Camadas" combinada com DDD, às vezes chamada de Clean/Hexagonal Architecture aplicada ao domínio.
Explicando as camadas:
- domain/: regras de negócio puras. Não importa Express, não importa Prisma. Só TypeScript puro.
- application/: orquestra o domínio para cumprir um caso de uso (um Command ou uma Query).
- infrastructure/: implementações concretas (banco de dados, fila, cache).
- interfaces/: a "porta de entrada" — neste caso HTTP (rotas Express), mas poderia ser CLI, gRPC, etc.
4. Passo 1 — Subir a infraestrutura com Docker Compose
Crie docker-compose.yml na raiz:
version: "3.8"
services:
postgres:
image: postgres:16
restart: always
environment:
POSTGRES_USER: admin
POSTGRES_PASSWORD: admin
POSTGRES_MULTIPLE_DATABASES: orders_db,inventory_db,query_db
ports:
- "5432:5432"
volumes:
- ./scripts/create-multiple-dbs.sh:/docker-entrypoint-initdb.d/create-multiple-dbs.sh
- pg_data:/var/lib/postgresql/data
rabbitmq:
image: rabbitmq:3.13-management
restart: always
ports:
- "5672:5672" # protocolo AMQP
- "15672:15672" # painel web (usuário/senha: guest/guest)
redis:
image: redis:7-alpine
restart: always
ports:
- "6379:6379"
volumes:
pg_data:
Crie o script que cria os 3 bancos automaticamente (scripts/create-multiple-dbs.sh):
#!/bin/bash
set -e
set -u
function create_database() {
local database=$1
echo "Criando banco '$database'"
psql -v ON_ERROR_STOP=1 --username "$POSTGRES_USER" <<-EOSQL
CREATE DATABASE $database;
EOSQL
}
if [ -n "$POSTGRES_MULTIPLE_DATABASES" ]; then
for db in $(echo $POSTGRES_MULTIPLE_DATABASES | tr ',' ' '); do
create_database $db
done
fi
chmod +x scripts/create-multiple-dbs.sh
docker compose up -d
Confirme que tudo subiu:
docker compose ps
Acesse o painel do RabbitMQ em http://localhost:15672 (guest/guest) — útil para "ver" as mensagens passando enquanto testa.
5. Passo 2 — Order Service (lado de escrita / Commands)
5.1 Setup inicial
mkdir order-service && cd order-service
npm init -y
npm install express amqplib @prisma/client zod
npm install -D typescript ts-node-dev @types/express @types/node @types/amqplib prisma
npx tsc --init
npx prisma init
tsconfig.json (ajuste principal — habilite decorators não é necessário aqui, mantemos simples):
{
"compilerOptions": {
"target": "ES2020",
"module": "CommonJS",
"rootDir": "src",
"outDir": "dist",
"strict": true,
"esModuleInterop": true,
"skipLibCheck": true,
"resolveJsonModule": true
}
}
.env:
DATABASE_URL="postgresql://admin:admin@localhost:5432/orders_db"
RABBITMQ_URL="amqp://guest:guest@localhost:5672"
PORT=3001
5.2 Prisma Schema
prisma/schema.prisma:
generator client {
provider = "prisma-client-js"
}
datasource db {
provider = "postgresql"
url = env("DATABASE_URL")
}
model Order {
id String @id @default(uuid())
customerId String
status String @default("PENDING") // PENDING | CONFIRMED | CANCELLED
totalCents Int
currency String @default("AOA")
createdAt DateTime @default(now())
items OrderItem[]
}
model OrderItem {
id String @id @default(uuid())
orderId String
order Order @relation(fields: [orderId], references: [id])
productId String
quantity Int
unitCents Int
}
npx prisma migrate dev --name init
Nota didática: o Prisma aqui vive só na camada
infrastructure/. O domínio (Order.ts) não conhece o Prisma — isso é o que permite trocar de banco de dados sem tocar nas regras de negócio.
5.3 Camada de Domínio
src/domain/value-objects/Money.ts:
export class Money {
private constructor(private readonly cents: number, private readonly currency: string) {
if (cents < 0) throw new Error("Valor monetário não pode ser negativo");
}
static fromCents(cents: number, currency = "AOA"): Money {
return new Money(cents, currency);
}
add(other: Money): Money {
if (other.currency !== this.currency) {
throw new Error("Não é possível somar moedas diferentes");
}
return new Money(this.cents + other.cents, this.currency);
}
get value(): number {
return this.cents;
}
get currencyCode(): string {
return this.currency;
}
}
src/domain/events/OrderCreated.ts:
export interface OrderCreatedPayload {
orderId: string;
customerId: string;
items: { productId: string; quantity: number }[];
totalCents: number;
currency: string;
occurredAt: string;
}
export class OrderCreated {
static eventName = "order.created";
constructor(public readonly payload: OrderCreatedPayload) {}
}
src/domain/entities/Order.ts — o Aggregate Root:
import { Money } from "../value-objects/Money";
import { OrderCreated } from "../events/OrderCreated";
export interface OrderItemProps {
productId: string;
quantity: number;
unitCents: number;
}
export class Order {
private domainEvents: OrderCreated[] = [];
private constructor(
public readonly id: string,
public readonly customerId: string,
public readonly items: OrderItemProps[],
public status: "PENDING" | "CONFIRMED" | "CANCELLED",
public readonly currency: string
) {}
static create(id: string, customerId: string, items: OrderItemProps[], currency = "AOA"): Order {
if (items.length === 0) {
throw new Error("Um pedido precisa ter ao menos um item");
}
const order = new Order(id, customerId, items, "PENDING", currency);
order.domainEvents.push(
new OrderCreated({
orderId: order.id,
customerId: order.customerId,
items: items.map((i) => ({ productId: i.productId, quantity: i.quantity })),
totalCents: order.total().value,
currency: order.currency,
occurredAt: new Date().toISOString(),
})
);
return order;
}
total(): Money {
return this.items.reduce(
(acc, item) => acc.add(Money.fromCents(item.unitCents * item.quantity, this.currency)),
Money.fromCents(0, this.currency)
);
}
pullDomainEvents(): OrderCreated[] {
const events = [...this.domainEvents];
this.domainEvents = [];
return events;
}
}
Repare:
Order.create(...)é onde a regra "um pedido precisa ter itens" mora — e não num controller ou numifsolto em algum lugar. Isso é DDD na prática: o domínio protege suas próprias invariantes.
src/domain/repositories/IOrderRepository.ts:
import { Order } from "../entities/Order";
export interface IOrderRepository {
save(order: Order): Promise<void>;
}
5.4 Camada de Aplicação (CQRS — Command side)
src/application/commands/CreateOrder.command.ts:
export interface CreateOrderCommand {
customerId: string;
items: { productId: string; quantity: number; unitCents: number }[];
}
src/application/commands/CreateOrder.handler.ts:
import { randomUUID } from "crypto";
import { Order } from "../../domain/entities/Order";
import { IOrderRepository } from "../../domain/repositories/IOrderRepository";
import { CreateOrderCommand } from "./CreateOrder.command";
import { IEventPublisher } from "../../domain/events/IEventPublisher";
export class CreateOrderHandler {
constructor(
private readonly repository: IOrderRepository,
private readonly publisher: IEventPublisher
) {}
async execute(command: CreateOrderCommand): Promise<string> {
const order = Order.create(randomUUID(), command.customerId, command.items);
await this.repository.save(order);
// Publica os eventos de domínio DEPOIS de persistir com sucesso
for (const event of order.pullDomainEvents()) {
await this.publisher.publish(event);
}
return order.id;
}
}
Crie também a interface src/domain/events/IEventPublisher.ts:
export interface IEventPublisher {
publish(event: { payload: unknown } & { constructor: { eventName?: string } }): Promise<void>;
}
Este handler é o caso de uso. Note que ele não sabe se o repositório usa Prisma nem se o publisher usa RabbitMQ — ele só depende de interfaces. Isso é o princípio de inversão de dependência, essencial para testar a lógica sem precisar subir banco de dados ou fila.
5.5 Camada de Infraestrutura
src/infrastructure/database/PrismaOrderRepository.ts:
import { PrismaClient } from "@prisma/client";
import { Order } from "../../domain/entities/Order";
import { IOrderRepository } from "../../domain/repositories/IOrderRepository";
export class PrismaOrderRepository implements IOrderRepository {
constructor(private readonly prisma: PrismaClient) {}
async save(order: Order): Promise<void> {
await this.prisma.order.create({
data: {
id: order.id,
customerId: order.customerId,
status: order.status,
totalCents: order.total().value,
currency: order.currency,
items: {
create: order.items.map((item) => ({
productId: item.productId,
quantity: item.quantity,
unitCents: item.unitCents,
})),
},
},
});
}
}
src/infrastructure/messaging/RabbitMQPublisher.ts:
import amqplib, { Channel, Connection } from "amqplib";
import { IEventPublisher } from "../../domain/events/IEventPublisher";
const EXCHANGE = "orders_exchange";
export class RabbitMQPublisher implements IEventPublisher {
private connection!: Connection;
private channel!: Channel;
async connect(url: string): Promise<void> {
this.connection = await amqplib.connect(url);
this.channel = await this.connection.createChannel();
await this.channel.assertExchange(EXCHANGE, "topic", { durable: true });
}
async publish(event: { payload: unknown; constructor: { eventName?: string } }): Promise<void> {
const routingKey = event.constructor.eventName ?? "unknown.event";
const buffer = Buffer.from(JSON.stringify(event.payload));
this.channel.publish(EXCHANGE, routingKey, buffer, { persistent: true });
console.log(`[order-service] evento publicado: ${routingKey}`);
}
}
Por que "topic exchange"? Um exchange do tipo
topicpermite várias filas assinarem padrões de routing key (ex:order.*). Isso deixa a arquitetura flexível: amanhã você pode ter mais serviços interessados em eventos de pedido sem alterar o Order Service.
5.6 Camada de Interface (HTTP)
src/interfaces/http/order.controller.ts:
import { Request, Response } from "express";
import { z } from "zod";
import { CreateOrderHandler } from "../../application/commands/CreateOrder.handler";
const createOrderSchema = z.object({
customerId: z.string().min(1),
items: z
.array(
z.object({
productId: z.string().min(1),
quantity: z.number().int().positive(),
unitCents: z.number().int().nonnegative(),
})
)
.min(1),
});
export class OrderController {
constructor(private readonly createOrderHandler: CreateOrderHandler) {}
create = async (req: Request, res: Response) => {
const parsed = createOrderSchema.safeParse(req.body);
if (!parsed.success) {
return res.status(400).json({ error: parsed.error.flatten() });
}
try {
const orderId = await this.createOrderHandler.execute(parsed.data);
return res.status(201).json({ orderId });
} catch (err) {
return res.status(400).json({ error: (err as Error).message });
}
};
}
src/interfaces/http/order.routes.ts:
import { Router } from "express";
import { OrderController } from "./order.controller";
export function orderRoutes(controller: OrderController): Router {
const router = Router();
router.post("/orders", controller.create);
return router;
}
5.7 Montando tudo — server.ts
import express from "express";
import { PrismaClient } from "@prisma/client";
import { PrismaOrderRepository } from "./infrastructure/database/PrismaOrderRepository";
import { RabbitMQPublisher } from "./infrastructure/messaging/RabbitMQPublisher";
import { CreateOrderHandler } from "./application/commands/CreateOrder.handler";
import { OrderController } from "./interfaces/http/order.controller";
import { orderRoutes } from "./interfaces/http/order.routes";
async function bootstrap() {
const prisma = new PrismaClient();
const publisher = new RabbitMQPublisher();
await publisher.connect(process.env.RABBITMQ_URL!);
const repository = new PrismaOrderRepository(prisma);
const createOrderHandler = new CreateOrderHandler(repository, publisher);
const controller = new OrderController(createOrderHandler);
const app = express();
app.use(express.json());
app.use(orderRoutes(controller));
const port = process.env.PORT || 3001;
app.listen(port, () => console.log(`[order-service] rodando na porta ${port}`));
}
bootstrap().catch((err) => {
console.error("Falha ao iniciar order-service", err);
process.exit(1);
});
Adicione o script no package.json:
"scripts": {
"dev": "ts-node-dev --respawn src/server.ts"
}
Teste:
npm run dev
curl -X POST http://localhost:3001/orders \
-H "Content-Type: application/json" \
-d '{"customerId":"cliente-1","items":[{"productId":"prod-1","quantity":2,"unitCents":5000}]}'
Você deve ver no terminal o log evento publicado: order.created e no painel do RabbitMQ (aba Exchanges) o exchange orders_exchange recebendo mensagens.
6. Passo 3 — Inventory Service (consumidor de eventos)
6.1 Setup
mkdir inventory-service && cd inventory-service
npm init -y
npm install express amqplib @prisma/client
npm install -D typescript ts-node-dev @types/express @types/node @types/amqplib prisma
npx tsc --init
npx prisma init
.env:
DATABASE_URL="postgresql://admin:admin@localhost:5432/inventory_db"
RABBITMQ_URL="amqp://guest:guest@localhost:5672"
PORT=3002
prisma/schema.prisma:
generator client {
provider = "prisma-client-js"
}
datasource db {
provider = "postgresql"
url = env("DATABASE_URL")
}
model Stock {
productId String @id
quantity Int @default(0)
}
model ProcessedEvent {
id String @id // usaremos um id do evento para idempotência
processedAt DateTime @default(now())
}
npx prisma migrate dev --name init
Por que a tabela
ProcessedEvent? Mensageria não garante "exactly-once" por padrão (o padrão realista é "at-least-once" — uma mensagem pode ser entregue mais de uma vez, ex: após reconexão). Para não abater o estoque duas vezes pelo mesmo pedido, registramos os eventos já processados e ignoramos repetições. Isso é o padrão de idempotência de consumidor.
6.2 Consumer
src/infrastructure/messaging/RabbitMQConsumer.ts:
import amqplib from "amqplib";
import { PrismaClient } from "@prisma/client";
import { randomUUID } from "crypto";
const EXCHANGE = "orders_exchange";
const QUEUE = "inventory_service.order_created";
const ROUTING_KEY = "order.created";
export async function startInventoryConsumer(rabbitUrl: string, prisma: PrismaClient) {
const connection = await amqplib.connect(rabbitUrl);
const channel = await connection.createChannel();
await channel.assertExchange(EXCHANGE, "topic", { durable: true });
await channel.assertQueue(QUEUE, { durable: true });
await channel.bindQueue(QUEUE, EXCHANGE, ROUTING_KEY);
console.log(`[inventory-service] aguardando mensagens em "${QUEUE}"...`);
channel.consume(QUEUE, async (msg) => {
if (!msg) return;
try {
const event = JSON.parse(msg.content.toString());
const eventId = msg.properties.messageId ?? randomUUID();
const alreadyProcessed = await prisma.processedEvent.findUnique({ where: { id: eventId } });
if (alreadyProcessed) {
channel.ack(msg);
return;
}
for (const item of event.items as { productId: string; quantity: number }[]) {
await prisma.stock.upsert({
where: { productId: item.productId },
update: { quantity: { decrement: item.quantity } },
create: { productId: item.productId, quantity: -item.quantity },
});
}
await prisma.processedEvent.create({ data: { id: eventId } });
console.log(`[inventory-service] estoque atualizado para o pedido ${event.orderId}`);
channel.ack(msg);
} catch (err) {
console.error("[inventory-service] erro ao processar evento", err);
// não reenfileira automaticamente para evitar loop infinito num erro persistente;
// em produção, isso normalmente vai para uma fila de erro (dead-letter queue)
channel.nack(msg, false, false);
}
});
}
src/server.ts:
import express from "express";
import { PrismaClient } from "@prisma/client";
import { startInventoryConsumer } from "./infrastructure/messaging/RabbitMQConsumer";
async function bootstrap() {
const prisma = new PrismaClient();
await startInventoryConsumer(process.env.RABBITMQ_URL!, prisma);
const app = express();
app.get("/stock/:productId", async (req, res) => {
const stock = await prisma.stock.findUnique({ where: { productId: req.params.productId } });
res.json(stock ?? { productId: req.params.productId, quantity: 0 });
});
const port = process.env.PORT || 3002;
app.listen(port, () => console.log(`[inventory-service] rodando na porta ${port}`));
}
bootstrap();
Teste: crie um pedido novamente via Order Service e depois confira:
curl http://localhost:3002/stock/prod-1
Nota conceitual: repare que
msg.properties.messageIdé usado como chave de idempotência — noRabbitMQPublisher.tsdo Order Service seria uma boa prática definir explicitamente essemessageId(ex: um UUID gerado por evento) ao publicar. Deixei como exercício para você aplicar, é uma ótima forma de fixar o conceito.
7. Passo 4 — Query Service (lado de leitura CQRS + Redis)
7.1 Setup
mkdir query-service && cd query-service
npm init -y
npm install express amqplib @prisma/client ioredis
npm install -D typescript ts-node-dev @types/express @types/node @types/amqplib prisma
npx tsc --init
npx prisma init
.env:
DATABASE_URL="postgresql://admin:admin@localhost:5432/query_db"
RABBITMQ_URL="amqp://guest:guest@localhost:5672"
REDIS_URL="redis://localhost:6379"
PORT=3003
prisma/schema.prisma — repare que este modelo é desnormalizado, pensado só para leitura rápida (não tem as mesmas normalizações do Order Service):
generator client {
provider = "prisma-client-js"
}
datasource db {
provider = "postgresql"
url = env("DATABASE_URL")
}
model OrderReadModel {
id String @id
customerId String
totalCents Int
currency String
itemsJson String // guardamos os itens já "achatados" em JSON para não precisar de JOIN
createdAt DateTime @default(now())
}
npx prisma migrate dev --name init
7.2 Consumer que alimenta o modelo de leitura
src/infrastructure/messaging/RabbitMQConsumer.ts:
import amqplib from "amqplib";
import { PrismaClient } from "@prisma/client";
import Redis from "ioredis";
const EXCHANGE = "orders_exchange";
const QUEUE = "query_service.order_created";
const ROUTING_KEY = "order.created";
export async function startQueryConsumer(rabbitUrl: string, prisma: PrismaClient, redis: Redis) {
const connection = await amqplib.connect(rabbitUrl);
const channel = await connection.createChannel();
await channel.assertExchange(EXCHANGE, "topic", { durable: true });
await channel.assertQueue(QUEUE, { durable: true });
await channel.bindQueue(QUEUE, EXCHANGE, ROUTING_KEY);
channel.consume(QUEUE, async (msg) => {
if (!msg) return;
try {
const event = JSON.parse(msg.content.toString());
await prisma.orderReadModel.upsert({
where: { id: event.orderId },
update: {},
create: {
id: event.orderId,
customerId: event.customerId,
totalCents: event.totalCents,
currency: event.currency,
itemsJson: JSON.stringify(event.items),
},
});
// Já deixamos o cache "quente" assim que o dado chega — cache-aside no write path
await redis.set(`order:${event.orderId}`, JSON.stringify(event), "EX", 60 * 5);
console.log(`[query-service] read model atualizado para o pedido ${event.orderId}`);
channel.ack(msg);
} catch (err) {
console.error("[query-service] erro ao processar evento", err);
channel.nack(msg, false, false);
}
});
}
7.3 Query Handler com cache Redis (padrão cache-aside)
src/infrastructure/cache/RedisClient.ts:
import Redis from "ioredis";
export function createRedisClient(url: string): Redis {
return new Redis(url);
}
src/application/queries/GetOrderById.handler.ts:
import { PrismaClient } from "@prisma/client";
import Redis from "ioredis";
const CACHE_TTL_SECONDS = 60 * 5;
export class GetOrderByIdHandler {
constructor(private readonly prisma: PrismaClient, private readonly redis: Redis) {}
async execute(orderId: string) {
const cacheKey = `order:${orderId}`;
// 1. Tenta o cache primeiro
const cached = await this.redis.get(cacheKey);
if (cached) {
console.log(`[query-service] cache HIT para ${orderId}`);
return JSON.parse(cached);
}
console.log(`[query-service] cache MISS para ${orderId}`);
// 2. Cache miss → busca no banco de leitura
const order = await this.prisma.orderReadModel.findUnique({ where: { id: orderId } });
if (!order) return null;
const result = {
orderId: order.id,
customerId: order.customerId,
totalCents: order.totalCents,
currency: order.currency,
items: JSON.parse(order.itemsJson),
};
// 3. Preenche o cache para as próximas leituras
await this.redis.set(cacheKey, JSON.stringify(result), "EX", CACHE_TTL_SECONDS);
return result;
}
}
Este é o padrão clássico cache-aside (também chamado lazy-loading): a aplicação verifica o cache, e só vai ao banco se não encontrar nada, populando o cache em seguida. Combinado com o passo anterior (o consumer já "esquenta" o cache assim que o evento chega), a maioria das leituras nunca toca o Postgres.
7.4 Rota HTTP e bootstrap
src/interfaces/http/order.routes.ts:
import { Router } from "express";
import { GetOrderByIdHandler } from "../../application/queries/GetOrderById.handler";
export function orderQueryRoutes(handler: GetOrderByIdHandler): Router {
const router = Router();
router.get("/orders/:id", async (req, res) => {
const order = await handler.execute(req.params.id);
if (!order) return res.status(404).json({ error: "Pedido não encontrado" });
res.json(order);
});
return router;
}
src/server.ts:
import express from "express";
import { PrismaClient } from "@prisma/client";
import { createRedisClient } from "./infrastructure/cache/RedisClient";
import { startQueryConsumer } from "./infrastructure/messaging/RabbitMQConsumer";
import { GetOrderByIdHandler } from "./application/queries/GetOrderById.handler";
import { orderQueryRoutes } from "./interfaces/http/order.routes";
async function bootstrap() {
const prisma = new PrismaClient();
const redis = createRedisClient(process.env.REDIS_URL!);
await startQueryConsumer(process.env.RABBITMQ_URL!, prisma, redis);
const handler = new GetOrderByIdHandler(prisma, redis);
const app = express();
app.use(express.json());
app.use(orderQueryRoutes(handler));
const port = process.env.PORT || 3003;
app.listen(port, () => console.log(`[query-service] rodando na porta ${port}`));
}
bootstrap();
8. Passo 5 — Testando o fluxo completo
Abra 3 terminais (um por serviço) e rode npm run dev em cada um. Depois:
# 1. Cria o pedido (Order Service)
curl -X POST http://localhost:3001/orders \
-H "Content-Type: application/json" \
-d '{"customerId":"cliente-1","items":[{"productId":"prod-1","quantity":3,"unitCents":2500}]}'
# resposta: {"orderId":"xxxxxxxx-xxxx-..."}
# 2. Confere que o Inventory Service abateu o estoque
curl http://localhost:3002/stock/prod-1
# 3. Confere que o Query Service já responde com o pedido (via Redis)
curl http://localhost:3003/orders/xxxxxxxx-xxxx-...
Se tudo funcionar, você acabou de ver na prática:
- Um Command (
POST /orders) alterando o estado através de um Aggregate DDD. - Um Domain Event (
OrderCreated) desacoplando o Order Service do resto do sistema. - Dois consumidores independentes (Inventory e Query) reagindo ao mesmo evento, cada um dentro do seu próprio Bounded Context.
- Uma Query (
GET /orders/:id) servida por um modelo de leitura desnormalizado e cacheado — o coração do CQRS.
9. Erros comuns ao estudar isso (e como evitar)
| Sintoma | Causa provável |
|---|---|
ECONNREFUSED ao conectar no RabbitMQ/Postgres/Redis |
Docker Compose não terminou de subir — espere alguns segundos ou rode docker compose logs -f
|
Query Service devolve 404 logo após criar o pedido |
Normal: é a consistência eventual. O evento ainda não chegou. Tente de novo em 1 segundo |
| Estoque abatido duas vezes para o mesmo pedido | Esqueceu de checar ProcessedEvent (idempotência) ou não reiniciou o serviço após adicionar essa lógica |
Cannot find module '@prisma/client' |
Esqueceu de rodar npx prisma generate (o migrate dev já faz isso, mas se só editou o schema, rode manualmente) |
10. Próximos passos (para quando quiser evoluir o projeto)
Este guia cobre o essencial para aprender os conceitos na prática. Se quiser aprofundar, aqui estão caminhos naturais de evolução:
- Outbox Pattern: hoje, entre "salvar o pedido no Postgres" e "publicar no RabbitMQ" existe uma pequena janela de falha (e se o processo cair entre os dois passos?). O padrão Transactional Outbox resolve isso salvando o evento na mesma transação do banco e tendo um processo separado que lê essa tabela e publica.
- Saga Pattern: para fluxos com múltiplos passos que precisam ser desfeitos em caso de erro (ex: pedido criado → pagamento falha → precisa devolver o estoque), pesquise sobre Sagas coreografadas (por eventos, como já fizemos aqui) ou orquestradas (um serviço central comanda o fluxo).
-
API Gateway: hoje o cliente precisa saber que existem 2 portas diferentes (
3001para escrever,3003para ler). Um Gateway (ex: com Express mesmo, ou algo como Kong/Traefik) unifica isso. -
Dead Letter Queue (DLQ): hoje, ao dar
nacksem reenfileirar, a mensagem problemática simplesmente desaparece. Configure uma fila de erro no RabbitMQ para investigar depois. -
Observabilidade: adicione logs estruturados (ex:
pino) e correlacione as requisições entre serviços com umtraceId(gerado no Order Service e propagado no payload do evento). -
Testes: como o
domain/e oapplication/não dependem de Express/Prisma/RabbitMQ, são fáceis de testar com testes unitários puros (ex: comvitestoujest), sem precisar subir Docker. - Autenticação entre serviços: hoje qualquer um pode chamar as rotas — em produção, pense em API Keys internas ou mTLS entre serviços.
11. Resumo mental (para revisar depois)
- DDD organiza o código pelo negócio: Entidades, Value Objects, Aggregates, Domain Events, Repositories.
- CQRS separa escrita (Commands, regras de negócio, consistência forte) de leitura (Queries, modelo desnormalizado, cache).
- RabbitMQ conecta os serviços de forma assíncrona via eventos, permitindo desacoplamento e resiliência.
- Redis, no lado de leitura, evita bater no banco toda hora — padrão cache-aside.
- Cada microserviço tem seu próprio banco de dados — isso é inegociável na arquitetura de microserviços.
- Sistemas distribuídos trocam consistência imediata por consistência eventual — e está tudo bem, desde que você desenhe o sistema sabendo disso.
Bons estudos! Qualquer parte que quiser que eu aprofunde (Outbox, Sagas, testes automatizados, etc.), é só pedir.
Top comments (0)