Introdução ao Apache Pulsar
Mensageria e streaming na mesma plataforma. Apache Pulsar na prática: do conceito ao laboratório em Python

Introdução ao Apache Pulsar
Mensageria e streaming na mesma plataforma. Apache Pulsar na prática: do conceito ao laboratório em Python
Introdução
Se você já trabalhou com sistemas distribuídos, provavelmente já ouviu falar do Apache Kafka e do RabbitMQ. Mas existem outras opções e aqui vamos falar sobre o Apache Pulsar.
O Pulsar é um daqueles projetos que chegou um pouco depois na festa, mas trouxe novidades suficientes para chamar atenção. Ele combina o melhor de dois mundos: o modelo de filas tradicional (estilo RabbitMQ) e o modelo de streaming de eventos (estilo Kafka), tudo em uma única plataforma.
Neste artigo você vai entender o que é o Pulsar, como ele funciona, por que ele tem chamado tanta atenção, e, principalmente, vai montar do zero um laboratório chamado Pulsar Lab: Pulsar no Docker, produtor e consumidor em Python, consumo compartilhado (Shared) e scripts para os problemas mais comuns de quem está começando.
Não é necessário clonar nenhum repositório: copie os arquivos abaixo, siga o passo a passo e tudo funcionará na sua máquina.
O que é Apache Pulsar?
O Apache Pulsar é uma plataforma de mensageria e streaming de eventos distribuída, criada pelo Yahoo! em 2013 e doada à Apache Software Foundation em 2016. Pense nele como um sistema de correios super organizado, que não só entrega cartas (mensagens) como também guarda tudo em um arquivo histórico (streaming), e ainda consegue atender várias empresas no mesmo cluster sem misturar as correspondências.
Algumas características fazem o Pulsar se destacar:
- Arquitetura separada: diferente do Kafka, no Pulsar quem processa as mensagens (brokers) é separado de quem armazena (bookies). Isso permite escalar cada parte de forma independente.
- Multi-tenancy nativo: vários times e aplicações podem compartilhar o mesmo cluster sem pisar no pé um do outro.
- Geo-replicação embutida: replicar dados entre data centers é configuração, não projeto.
- Suporte a filas e streaming: no mesmo tópico, você pode ter consumidores em modo fila (dividindo o trabalho) e em modo streaming (cada um lendo tudo, conforme o tipo de subscription).
Resumindo: o Pulsar é aquele canivete suíço que você não sabia que precisava até começar a usar.
Para que serve o Apache Pulsar?
O Pulsar é versátil e cobre vários cenários comuns:

Por que considerar o Pulsar?
- Escalabilidade real: brokers e storage escalam de forma independente.
- Modelo flexível de consumo: Exclusive, Shared, Failover e Key_Shared cobrem quase qualquer padrão.
- Schemas nativos: Avro, JSON e Protobuf com validação no broker.
- Tiered Storage: mensagens antigas podem ir para S3/GCS automaticamente.
- Pulsar Functions: processamento serverless leve embutido.
A curva de aprendizado é um pouco maior que a do RabbitMQ, mas você ganha uma plataforma que cresce com o sistema sem trocar de ferramenta no meio do caminho.
O projeto: Pulsar Lab
Vamos criar a pasta pulsar-lab com esta estrutura:
pulsar-lab/
├── docker-compose.yml # Pulsar 3.2.0 standalone
├── requirements.txt # Cliente Python
├── config.py # URLs e nome do tópico
├── pulsar_wait.py # Espera o broker ficar pronto
├── producer.py # Publica mensagens
├── consumer.py # Consome (Exclusive)
├── consumer_shared.py # Consome (Shared, vários workers)
├── scripts/
│ ├── wait-pulsar.sh # Espera via curl (opcional)
│ └── reset-subscriptions.sh # Libera subscriptions presas
└── .gitignore
Conceitos que vamos usar:

Pré-requisitos
- Docker e Docker Compose v2
- Python 3.10 ou superior
curl(para o script opcional de espera)
Crie e entre na pasta do projeto:
mkdir pulsar-lab && cd pulsar-lab
Passo 1: Subir o Pulsar com Docker Compose
Crie o arquivo docker-compose.yml:
services:
pulsar:
image: apachepulsar/pulsar:3.2.0
container_name: pulsar
command: bin/pulsar standalone
ports:
- "6650:6650"
- "8080:8080"
Suba o cluster e verificar se ele subiu:
docker compose up -d
docker compose ps



Para parar:
docker compose down
Importante: no modo standalone, o Pulsar leva cerca de 30 a 60 segundos após o
uppara criar o namespacepublic/default. Se os scripts Python rodarem antes disso, você veráTopicNotFound: Namespace not found. O módulopulsar_wait.py(abaixo) resolve isso automaticamente.
Passo 2: Ambiente Python
Crie requirements.txt:
pulsar-client==3.5.0
Crie o ambiente virtual e instale:
python3 -m venv .venv
source .venv/bin/activate # Windows: .venv\Scripts\activate
pip install -r requirements.txt

Crie .gitignore:
.venv/
__pycache__/
*.pyc
.pytest_cache/
Passo 3: Configuração compartilhada
Crie config.py:
"""Configuração compartilhada dos exemplos."""
PULSAR_SERVICE_URL = "pulsar://localhost:6650"
PULSAR_ADMIN_URL = "http://localhost:8080"
PULSAR_NAMESPACE = "public/default"
TOPIC = "minha-fila"
O nome curto minha-fila é expandido pelo broker para persistent://public/default/minha-fila no cluster standalone.
Passo 4: Esperar o broker ficar pronto
Crie pulsar_wait.py:
"""Aguarda o broker Pulsar e o namespace public/default ficarem prontos."""
from __future__ import annotations
import json
import time
import urllib.error
import urllib.request
from config import PULSAR_ADMIN_URL, PULSAR_NAMESPACE
_READY_NAMESPACE = PULSAR_NAMESPACE
def wait_for_pulsar(timeout: float = 120, interval: float = 2) -> None:
"""Bloqueia até a API admin responder com o namespace default."""
deadline = time.monotonic() + timeout
url = f"{PULSAR_ADMIN_URL}/admin/v2/namespaces/public"
while time.monotonic() < deadline:
try:
with urllib.request.urlopen(url, timeout=3) as response:
namespaces = json.loads(response.read().decode())
if _READY_NAMESPACE in namespaces:
return
except (urllib.error.URLError, urllib.error.HTTPError, OSError, TimeoutError, json.JSONDecodeError):
pass
print(" [*] Aguardando Pulsar inicializar (namespace public/default)...")
time.sleep(interval)
raise TimeoutError(
f"Pulsar não ficou pronto em {timeout:.0f}s. "
"Confira: docker compose ps && docker logs pulsar"
)
Opcional — script shell scripts/wait-pulsar.sh (crie a pasta scripts antes):
#!/usr/bin/env bash
set -euo pipefail
ADMIN_URL="${PULSAR_ADMIN_URL:-http://localhost:8080}"
TIMEOUT="${PULSAR_WAIT_TIMEOUT:-120}"
INTERVAL=2
DEADLINE=$((SECONDS + TIMEOUT))
echo "Aguardando Pulsar em ${ADMIN_URL}..."
while (( SECONDS < DEADLINE )); do
if curl -sf "${ADMIN_URL}/admin/v2/namespaces/public" 2>/dev/null | grep -q 'public/default'; then
echo "Pulsar pronto."
exit 0
fi
sleep "${INTERVAL}"
done
echo "Timeout: Pulsar não respondeu em ${TIMEOUT}s." >&2
echo "Verifique: docker compose ps && docker logs pulsar" >&2
exit 1
Torne executável: chmod +x scripts/wait-pulsar.sh
Passo 5: Produtor
Crie producer.py:
import sys
import pulsar
from config import PULSAR_SERVICE_URL, TOPIC
from pulsar_wait import wait_for_pulsar
def main() -> None:
messages = sys.argv[1:] or ["Olá, Pulsar!"]
wait_for_pulsar()
client = pulsar.Client(PULSAR_SERVICE_URL)
producer = client.create_producer(TOPIC)
for text in messages:
producer.send(text.encode("utf-8"))
print(f" [x] Mensagem enviada: {text}")
producer.close()
client.close()
if __name__ == "__main__":
main()
Sem argumentos, envia "Olá, Pulsar!". Com argumentos, envia cada um como mensagem:
python producer.py "msg 1" "msg 2" "msg 3"

Passo 6: Consumidor (Exclusive)
No Pulsar, o tipo padrão de subscription é Exclusive: apenas um consumidor pode estar conectado por vez na mesma subscription.
Crie consumer.py:
import sys
import pulsar
from config import PULSAR_SERVICE_URL, TOPIC
from pulsar_wait import wait_for_pulsar
_SUBSCRIPTION = "meu-app"
def main() -> None:
wait_for_pulsar()
client = pulsar.Client(PULSAR_SERVICE_URL)
try:
consumer = client.subscribe(TOPIC, subscription_name=_SUBSCRIPTION)
except pulsar.ConsumerBusy:
client.close()
print(
" [!] ConsumerBusy: já existe um consumidor Exclusive em 'meu-app'.\n"
" • Feche o outro terminal com consumer.py (Ctrl+C), ou\n"
" • Rode: ./scripts/reset-subscriptions.sh --force",
file=sys.stderr,
)
sys.exit(1)
print(" [*] Aguardando mensagens... (Ctrl+C para sair)")
try:
while True:
msg = consumer.receive()
try:
body = msg.data().decode("utf-8")
print(f" [x] Mensagem recebida: {body}")
consumer.acknowledge(msg)
except Exception:
consumer.negative_acknowledge(msg)
except KeyboardInterrupt:
print("\n [*] Encerrando consumidor...")
finally:
consumer.close()
client.close()
if __name__ == "__main__":
main()
acknowledge(msg)confirma o processamento com sucesso.negative_acknowledge(msg)pede reentrega em caso de falha.

Passo 7: Rodar o exemplo básico
Com o Pulsar no ar e o venv ativado:
Terminal 1 — consumidor (rode primeiro):
python consumer.py
Aguarde a linha:
[*] Aguardando mensagens... (Ctrl+C para sair)
Terminal 2 — produtor:
python producer.py
Saída esperada no consumidor:
[x] Mensagem recebida: Olá, Pulsar!
No produtor:
[x] Mensagem enviada: Olá, Pulsar!
Os logs
INFOdo cliente Pulsar são normais. Para reduzi-los, exportePULSAR_LOG_LEVEL=warnantes de rodar os scripts.
Passo 8: Bônus — vários consumidores (Shared)
Para ver o modelo de fila distribuída, use subscription Shared: o Pulsar divide as mensagens entre vários consumidores da mesma subscription.
Crie consumer_shared.py:
import os
import sys
import pulsar
from config import PULSAR_SERVICE_URL, TOPIC
from pulsar_wait import wait_for_pulsar
def main() -> None:
worker_id = os.environ.get("WORKER_ID") or (sys.argv[1] if len(sys.argv) > 1 else "worker-1")
wait_for_pulsar()
client = pulsar.Client(PULSAR_SERVICE_URL)
consumer = client.subscribe(
TOPIC,
subscription_name="workers",
consumer_type=pulsar.ConsumerType.Shared,
)
print(f" [{worker_id}] Aguardando mensagens (Shared)... (Ctrl+C para sair)")
try:
while True:
msg = consumer.receive()
try:
body = msg.data().decode("utf-8")
print(f" [{worker_id}] Mensagem recebida: {body}")
consumer.acknowledge(msg)
except Exception:
consumer.negative_acknowledge(msg)
except KeyboardInterrupt:
print(f"\n [{worker_id}] Encerrando consumidor...")
finally:
consumer.close()
client.close()
if __name__ == "__main__":
main()
Terminal 1:
WORKER_ID=worker-1 python consumer_shared.py
Terminal 2:
WORKER_ID=worker-2 python consumer_shared.py
Terminal 3:
python producer.py "a" "b" "c" "d" "e" "f"
Você verá as seis mensagens distribuídas entre worker-1 e worker-2 — sem configuração extra no broker.
Tipos de subscription (referência)

Solução de problemas
TopicNotFound: Namespace not found
O broker ainda não criou public/default. Os scripts chamam wait_for_pulsar() e exibem:
[*] Aguardando Pulsar inicializar (namespace public/default)...
Se o timeout estourar:
docker compose ps
docker logs pulsar
Confirme manualmente:
curl -s http://localhost:8080/admin/v2/namespaces/public
# Deve incluir "public/default"
ConsumerBusy: Exclusive consumer is already connected
Outro consumer.py já está conectado à subscription meu-app.
- Feche o terminal antigo com Ctrl+C, ou
- Liste processos:
pgrep -af consumer.py - Encerre:
kill <PID> - Se ainda falhar, use o script de reset (abaixo) com
--force
Crie scripts/reset-subscriptions.sh:
#!/usr/bin/env bash
# Desconecta consumidores presos nas subscriptions do lab.
set -euo pipefail
CONTAINER="${PULSAR_CONTAINER:-pulsar}"
TOPIC="persistent://public/default/minha-fila"
FORCE=false
if [[ "${1:-}" == "--force" || "${1:-}" == "-f" ]]; then
FORCE=true
fi
unsubscribe() {
local sub="$1"
local args=(-s "${sub}")
if [[ "${FORCE}" == true ]]; then
args=(-f "${args[@]}")
fi
docker exec "${CONTAINER}" bin/pulsar-admin topics unsubscribe "${TOPIC}" "${args[@]}"
}
for sub in meu-app workers; do
if docker exec "${CONTAINER}" bin/pulsar-admin topics list-subscriptions "${TOPIC}" 2>/dev/null | grep -qx "${sub}"; then
echo "Desconectando consumidores da subscription '${sub}'..."
unsubscribe "${sub}" || {
echo " Falhou (consumidor ainda ativo?). Tente: $0 --force" >&2
exit 1
}
fi
done
echo "Pronto. Rode python consumer.py novamente."
chmod +x scripts/reset-subscriptions.sh
./scripts/reset-subscriptions.sh --force
Checklist rápido
# 1. Infra
mkdir pulsar-lab && cd pulsar-lab
# (criar todos os arquivos deste artigo)
docker compose up -d
# 2. Python
python3 -m venv .venv && source .venv/bin/activate
pip install -r requirements.txt
# 3. Teste
# Terminal A
python consumer.py
# Terminal B
python producer.py
E agora?
Com o Pulsar Lab funcionando, os próximos passos naturais são:
- Schemas (Avro, JSON, Protobuf) para contratos de dados rígidos
- Pulsar Functions para transformações serverless no broker
- Geo-replicação para sistemas multi-região
- Tiered storage para reduzir custo de retenção longa
Explore a UI em http://localhost:8080, quebre o lab, conserte de novo — é assim que a gente aprende.
Conclusão
O Apache Pulsar une mensageria tradicional e streaming de eventos em uma plataforma que escala de verdade. Neste artigo você montou o Pulsar Lab completo: Docker Compose, espera automática do namespace, produtor, consumidor Exclusive, workers Shared e scripts para os erros mais comuns do dia a dia.
Se o seu sistema está crescendo, se você está cansado de manter ferramentas separadas para fila e streaming, ou se geo-replicação é requisito, vale dar uma chance ao Pulsar. O laboratório que você acabou de construir é um ótimo ponto de partida — sem depender de repositório externo algum.
Até a próxima!
메타데이터
- post_id
- 58f6997cfdb9
- slug
- introdução-ao-apache-pulsar-58f6997cfdb9
- url
- https://medium.com/@habbema/introdu%C3%A7%C3%A3o-ao-apache-pulsar-58f6997cfdb9
- canonical_url
- https://medium.com/@habbema/introdu%C3%A7%C3%A3o-ao-apache-pulsar-58f6997cfdb9
- author_url
- https://medium.com/@habbema
- status
- ok
- fetched_at
- 2026-07-17 15:39:32