← Back to list

Introdução ao Apache Pulsar

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

Hugo Habbema · 2026-05-24 15:46 · 5 claps · 8.0 min read
#apache-pulsar
Open on Medium ↗
Wiki topics: 🎬 · Film & Television 📊 · Economic Policy

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 up para criar o namespace public/default. Se os scripts Python rodarem antes disso, você verá TopicNotFound: Namespace not found. O módulo pulsar_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 INFO do cliente Pulsar são normais. Para reduzi-los, exporte PULSAR_LOG_LEVEL=warn antes 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.

  1. Feche o terminal antigo com Ctrl+C, ou
  2. Liste processos: pgrep -af consumer.py
  3. Encerre: kill <PID>
  4. 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