Become a member!

PostgreSQL como job queue: você precisa mesmo de Redis ou RabbitMQ?

🌐
Este artigo também está disponível em outros idiomas:
🇪🇸 Español  •  🇩🇪 Deutsch  •  🇬🇧 English

Introdução

Neste artigo vamos ver por que faz sentido usar o Postgres como fila, analisar o design do schema, implementar um polling de tarefas seguro na concorrência com SELECT … FOR UPDATE SKIP LOCKED e garantir que um único worker processe uma determinada tarefa a cada momento. No final você terá o roteiro de um sistema confiável e transacionalmente seguro para processar tarefas em background, tudo rodando no Postgres.

Mar aberto sob um céu nublado ao pôr do sol, foto de capa de Daniele Teti


Sumário

  1. Por que usar o Postgres como job queue?
  2. Conceitos-chave e definições
  3. Design do schema do banco
  4. Locking no Postgres: FOR UPDATE SKIP LOCKED
  5. Inserindo novas tarefas
  6. Polling das tarefas a processar
  7. Atualizando o status da tarefa
  8. Workers paralelos e questões de concorrência
  9. Retentativas e tratamento de erros
  10. Prós e contras
  11. Exemplo de implementação passo a passo
  12. Dicas práticas de uso
  13. Conclusão e uma pergunta para refletir

Os assinantes do PATREON têm acesso a um projeto Delphi completo que implementa um sistema de job queue.

No artigo do PATREON você encontra 2 projetos:

  • JobProducer.dproj: é o programa cliente que pede a execução de algum job. No exemplo ele pode pedir 2 tipos de job (“send_email”, “create_report”), mas você pode estendê-lo para o que precisar.
  • QueueWorker.dproj: é o worker propriamente dito. Você pode abrir várias instâncias e verificar que o sistema atribui cada job a uma única instância de worker. Se o worker travar depois de pegar um job, esse job volta a ficar disponível para os outros workers.

1. Por que usar o Postgres como job queue?

Quando o assunto é job queue, os desenvolvedores costumam recorrer a ferramentas como Redis, RabbitMQ ou sistemas especializados como o Sidekiq (no ecossistema Ruby). Todas são opções perfeitamente válidas. Mas existem bons motivos para considerar o Postgres para as suas filas:

  1. Simplicidade: se a sua infraestrutura já depende do Postgres, você não precisa instalar nem configurar outro serviço. Isso reduz o custo de manutenção e o tempo para colocar a fila no ar.
  2. Consistência dos dados: o Postgres é um banco robusto e ACID. Gerenciar as tarefas dentro do mesmo banco garante consistência forte dos dados e integridade transacional.
  3. Confiabilidade: se você já confia no Postgres para guardar os dados críticos da aplicação, pode confiar nele também para gerenciar os jobs em background. Você evita criar vários pontos de falha, já que está tudo no mesmo lugar.
  4. Backup simples: você faz o backup da fila junto com os dados da aplicação em um único processo. Não precisa implementar nem manter uma estratégia de backup separada para mais um sistema.

Porém, é bom saber que o Postgres pode não ser a solução de fila perfeita para cenários de escala enorme, com dezenas de milhares de jobs publicados por segundo. Ele dá conta de muita coisa, mas você deve sempre medir o desempenho e otimizar de acordo, ou talvez considerar sistemas mais especializados nos casos extremos. Para a maioria das cargas pequenas e médias, o Postgres se sai muito bem.

🔔 Se você precisa de uma job queue com ainda mais recursos, porém mais simples de usar e com uma API REST limpa e simples, dê uma olhada no DMSContainer com o seu módulo Event Stream. Ele está pronto para uso e, com o DMSContainer, você começa a usar filas, enviar e-mails, gerar relatórios Excel e PDF (e muito mais) em minutos. Confira!


2. Conceitos-chave e definições

Antes de entrar nos detalhes, vamos definir alguns termos que vamos usar:

  • Job/tarefa: usamos os termos job e tarefa como sinônimos para indicar uma unidade de trabalho que precisa ser processada.
  • Producer: uma entidade (normalmente parte da sua aplicação) que cria as tarefas e as insere na fila.
  • Consumer/worker: um processo em background que pega as tarefas pendentes na fila, processa essas tarefas e atualiza o status delas.
  • FOR UPDATE SKIP LOCKED: um recurso do Postgres que permite travar as linhas enquanto você as seleciona, pulando as linhas já travadas por outra transação. É a peça-chave para evitar conflitos quando vários workers fazem polling das tarefas.

3. Design do schema do banco

Projetar uma tabela para guardar as tarefas é relativamente simples. Veja um schema mínimo que você pode usar:

CREATE TABLE tasks (
    id SERIAL PRIMARY KEY,
    status VARCHAR(50) NOT NULL DEFAULT 'pending',
    payload JSONB NOT NULL,
    created_at TIMESTAMP WITH TIME ZONE DEFAULT now(),
    updated_at TIMESTAMP WITH TIME ZONE DEFAULT now()
);

Veja o que faz cada coluna:

  • id: um identificador único para cada tarefa, gerado automaticamente.
  • status: indica o estado da tarefa, como ‘pending’, ‘in_progress’, ‘completed’ ou ‘failed’.
  • payload: contém os dados necessários para executar a tarefa. O JSONB dá flexibilidade para guardar estruturas de dados variadas.
  • created_at: timestamp que marca quando a tarefa foi criada.
  • updated_at: timestamp atualizado sempre que o status da tarefa muda.

Vamos manter este schema enxuto. Na prática você pode adicionar índices em campos JSON específicos, ou outras colunas como um nível de prioridade ou referências a outras tabelas.


4. Locking no Postgres: FOR UPDATE SKIP LOCKED

A estrela do show é o locking por linha do Postgres, mais precisamente o SELECT … FOR UPDATE SKIP LOCKED. Introduzido no PostgreSQL 9.5, o SKIP LOCKED garante que, se várias transações tentam travar as mesmas linhas, uma transação as trava enquanto as outras pulam as linhas que já estão travadas.

Isso significa que você pode rodar a mesma instrução SELECT em paralelo em vários workers sem risco de duplicação. Cada worker pega uma linha (ou seja, uma tarefa) e a trava para que ninguém mais possa pegá-la. É perfeito para distribuir tarefas entre vários workers, porque evita colisões.

Em resumo:

  • FOR UPDATE: obtém um lock sobre a(s) linha(s) para que nenhuma outra transação possa modificá-las até que o lock seja liberado.
  • SKIP LOCKED: diz ao Postgres para não esperar pelas linhas travadas, mas pulá-las e travar só as linhas que não estão travadas no momento.

5. Inserindo novas tarefas

Quando a sua aplicação cria uma nova tarefa, basta inserir uma linha na tabela tasks. Um exemplo simples:

INSERT INTO tasks (payload) 
VALUES ('{"task_type": "send_email", "recipient": "[email protected]"}');

O status assume o valor default ‘pending’, created_at assume o default now() e updated_at também recebe now(). É o DML padrão que você já conhece, sem surpresas.


6. Polling das tarefas a processar

Agora vem a parte importante: como buscar as tarefas de modo que um único worker processe uma determinada tarefa a cada vez? Suponha que temos vários processos (ou threads) worker, e cada um executa uma instrução SELECT para encontrar a próxima tarefa pendente.

Aqui vai uma versão simplificada:

BEGIN;

WITH cte_task AS (
    SELECT id
    FROM tasks
    WHERE status = 'pending'
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 1
)
UPDATE tasks
SET status = 'in_progress',
    updated_at = now()
FROM cte_task
WHERE tasks.id = cte_task.id
RETURNING tasks.id, tasks.payload;

Fazemos tudo isso em uma única transação:

  1. WITH cte_task: esta common table expression trava exatamente uma tarefa em status ‘pending’. Usamos ORDER BY id só para ter uma ordenação determinística, mas você pode ordenar pela data de criação ou pela prioridade, se preferir.
  2. FOR UPDATE SKIP LOCKED: trava a linha e a pula se outro worker já tiver travado a mesma linha.
  3. LIMIT 1: garante que travamos uma tarefa de cada vez.
  4. UPDATE: depois que a CTE encontra a linha, mudamos o status dela para in_progress, atribuindo de fato a tarefa a este worker.
  5. RETURNING: recebemos de volta o id e o payload da tarefa em que vamos trabalhar.

Se a query acima retornar zero linhas, não há tarefas em estado ‘pending’. O worker pode simplesmente fazer commit e dormir um pouco antes de tentar de novo.

Se retornar uma linha, o worker pode seguir e processar a tarefa. Terminado o trabalho, o worker pode executar outro UPDATE para marcar a tarefa como ‘completed’, ou talvez ‘failed’ se o processamento deu erro.

Essa abordagem garante que cada tarefa seja travada e atualizada em um único passo atômico, impedindo que outros workers peguem a mesma tarefa ao mesmo tempo.


7. Atualizando o status da tarefa

Quando um worker termina a tarefa, ele precisa atualizar a tabela tasks:

UPDATE tasks
SET status = 'completed',
    updated_at = now()
WHERE id = <task_id>;

Ou, se o processamento da tarefa falhar, você pode fazer:

UPDATE tasks
SET status = 'failed',
    updated_at = now()
WHERE id = <task_id>;

Talvez você também queira registrar as mensagens de erro, o número de retentativas ou outros metadados relevantes. Considere adicionar esses campos ao schema se você prevê que as suas tarefas podem falhar e precisar ser reprocessadas.


8. Workers paralelos e questões de concorrência

Usar FOR UPDATE SKIP LOCKED significa que vários workers podem rodar a mesma query de polling ao mesmo tempo sem conflito. Cada worker acaba travando linhas diferentes. Se um worker trava uma linha, os outros a pulam. Essa abordagem de concorrência é ideal para distribuir as tarefas entre um pool de workers.

Mas você precisa controlar com que frequência cada worker faz o polling da tabela. Um polling frequente demais pode gerar muita carga no banco, principalmente se você tiver muitos processos worker. É preciso equilibrar a rapidez em pegar as tarefas com a necessidade de não martelar o banco com SELECT demais.


9. Retentativas e tratamento de erros

No mundo real, o processamento de jobs nunca é 100% livre de erros. Você vai precisar de uma estratégia para as tarefas que falham. Suponha que um worker pegue uma tarefa, mas a execução do job falhe por um erro de rede ou algum outro problema. Você pode:

  1. Marcar a tarefa como ‘failed’ e registrar os detalhes do erro. Depois, outro subsistema ou um script específico pode procurar as tarefas ‘failed’ e decidir se as tenta de novo ou não.
  2. Tentar de novo automaticamente um número limitado de vezes. Você pode controlar uma coluna retry_count na tabela tasks. Se estiver abaixo de um certo limite, a tarefa volta ao status ‘pending’ e retry_count é incrementado.

Aqui vai um exemplo que marca uma tarefa como falha e incrementa um retry_count:

ALTER TABLE tasks ADD COLUMN retry_count INT NOT NULL DEFAULT 0;

BEGIN;

UPDATE tasks
SET status = 'failed',
    retry_count = retry_count + 1,
    updated_at = now()
WHERE id = <task_id>;

COMMIT;

A partir daí, você pode ter um processo separado ou um cron job que procura as tarefas failed com retry_count < 5 e as devolve para ‘pending’, recolocando-as na fila para mais uma tentativa. Em cenários mais simples, o próprio worker pode fazer retry_count = retry_count + 1 e parar de processar a tarefa depois de um número definido de tentativas.


10. Prós e contras

Por mais prático que seja usar o Postgres como job queue, não é uma solução universal. Abaixo está uma lista de prós e contras que reflete a opinião geral dos desenvolvedores:

Prós

  • Nenhuma infraestrutura extra: você já tem o Postgres, então não há mais nada para instalar ou manter.
  • Segurança transacional: se o seu código está fortemente ligado aos dados guardados no Postgres, ter as tarefas no mesmo lugar simplifica a consistência.
  • Locking por linha: o Postgres tem mecanismos de concorrência robustos. O FOR UPDATE SKIP LOCKED é bem testado e confiável.
  • Backup e restore: todos os dados, tarefas incluídas, podem ir para o mesmo backup.

Contras

  • Escalabilidade: se você tem um throughput altíssimo ou precisa de recursos avançados de fila (roteamento avançado, filas com prioridade etc.), sistemas especializados podem ser a melhor escolha.
  • Carga no banco: o polling das tarefas pode adicionar carga ao seu banco principal. Dá para mitigar com uma arquitetura bem pensada ou uma réplica, mas é mais uma coisa para considerar.
  • Recursos limitados: sistemas como RabbitMQ ou Kafka oferecem roteamento sofisticado de mensagens, fan-out e mais. O Postgres é mais simples nesse sentido, mas também menos flexível para certos padrões.

11. Exemplo de implementação passo a passo

Vamos percorrer um cenário hipotético com uma pequena startup de tecnologia: a Acme Email Services. Ela cuida de e-mails transacionais e precisa enfileirar as tarefas de envio de e-mail para não sobrecarregar o servidor SMTP.

  1. Criar a tabela tasks:

    CREATE TABLE tasks (
        id SERIAL PRIMARY KEY,
        status VARCHAR(50) NOT NULL DEFAULT 'pending',
        payload JSONB NOT NULL,
        created_at TIMESTAMP WITH TIME ZONE DEFAULT now(),
        updated_at TIMESTAMP WITH TIME ZONE DEFAULT now()
    );
    
  2. Inserir as tarefas (lado do producer, por exemplo a partir de um endpoint da API):

    INSERT INTO tasks (payload)
    VALUES ('{"task_type": "send_email", "subject": "Welcome!", "recipient": "[email protected]", "body": "Hello and welcome!"}');
    
  3. Processo worker (pseudocódigo em Python, por exemplo):

    import psycopg2
    import time
    
    def poll_and_process_tasks():
        while True:
            conn = psycopg2.connect("dbname=acme user=postgres password=postgres")
            conn.autocommit = False
            try:
                with conn.cursor() as cur:
                    cur.execute("""
                        WITH cte_task AS (
                            SELECT id, payload
                            FROM tasks
                            WHERE status = 'pending'
                            ORDER BY id
                            FOR UPDATE SKIP LOCKED
                            LIMIT 1
                        )
                        UPDATE tasks
                        SET status = 'in_progress',
                            updated_at = now()
                        FROM cte_task
                        WHERE tasks.id = cte_task.id
                        RETURNING tasks.id, tasks.payload;
                    """)
                    row = cur.fetchone()
                    if row:
                        task_id, task_payload = row
                        # Processa a tarefa
                        process_email(task_payload)  # a sua função personalizada
    
                        # Marca como concluída
                        cur.execute("""
                            UPDATE tasks
                            SET status = 'completed', updated_at = now()
                            WHERE id = %s
                        """, (task_id,))
                        conn.commit()
                    else:
                        conn.rollback()
                        # Nenhuma tarefa, dorme um pouco antes de verificar de novo
                        time.sleep(5)
            finally:
                conn.close()
    
    def process_email(payload):
        # Pseudocódigo para enviar o e-mail
        # payload['task_type'] == 'send_email'
        # Use payload['recipient'], payload['subject'] etc.
        print(f"Sending email to {payload['recipient']}...")
        # Aqui você envia o e-mail de verdade
    

    No script acima:

    • Conectamos ao banco e iniciamos uma transação.
    • Tentamos pegar uma única tarefa pending usando a técnica do FOR UPDATE SKIP LOCKED.
    • Se conseguimos uma, marcamos na hora como in_progress.
    • Em seguida processamos o e-mail (o código real de envio foi omitido para simplificar).
    • Se der certo, mudamos o status para completed.
    • Se não houver tarefas pendentes, fazemos rollback (para liberar locks ou atualizações parciais) e dormimos um pouco.

⭐ A versão Delphi completa do sistema de job queue será liberada para os apoiadores do PATREON nos próximos dias.

  1. Tratamento de erros: se acontecer um erro no envio do e-mail, você pode capturar a exceção, marcar a tarefa como failed e, se quiser, registrar o erro. Depois pode ter uma estratégia para recolocá-la na fila, se necessário.

12. Dicas práticas de uso

  • Indexação: se o volume de tarefas for muito alto, considere criar um índice em (status) ou (status, id) para acelerar a busca das tarefas pendentes.
  • Limite a frequência do polling: implemente um pequeno atraso nos workers ou uma estratégia de backoff para reduzir a carga no banco.
  • Use uma tabela separada: se a sua aplicação tem vários tipos de tarefa ou um número enorme de tarefas, você pode separá-las em tabelas diferentes ou até em schemas diferentes. Isso ajuda na organização e no ajuste de desempenho.
  • Monitoramento: fique de olho no número de tarefas em cada status. Ferramentas como o Grafana ou métricas personalizadas podem avisar quando começa a se formar um acúmulo (por exemplo, se as tarefas ficam ‘pending’ por tempo demais).
  • Tamanho da transação: preste atenção aos limites das transações. Se você tentar travar e processar centenas de tarefas em uma única transação, pode segurar os locks por tempo demais. Em geral você quer processar as tarefas uma a uma ou em lotes pequenos.

13. Conclusão e uma pergunta para refletir

A esta altura você já viu que usar o Postgres como job queue pode ser prático e elegante, principalmente se o seu ambiente já usa o Postgres para outras coisas. Você evita a complexidade de sistemas adicionais, ganha garantias transacionais sólidas e aproveita os recursos de concorrência do Postgres para distribuir as tarefas entre vários workers sem duplicação.

Com o PostgreSQL você pode montar um cluster de processos worker que dá conta de milhares de tarefas por dia de forma confiável. Guarde todos os logs e os metadados dos jobs no Postgres: fica simples consultar o histórico, analisar o desempenho e gerenciar a recuperação de erros. Soluções como DMSContainer/EventStream!, Redis, RabbitMQ ou Kafka podem ser mais adequadas em certos cenários especializados ou de grande escala, mas a abordagem da job queue no Postgres é uma forte candidata para muitas necessidades de pequena e média escala.

Pergunta para refletir: considerando a carga de trabalho e a infraestrutura da sua organização, você precisa dos recursos especializados de um sistema de filas dedicado, ou o Postgres consegue processar os seus jobs bem o bastante para simplificar a sua arquitetura? Mande um e-mail para consultoria e desenvolvimento especializados.

Artigos relacionados sobre PostgreSQL:


Referências


E essa é a história. Seja um projeto paralelo, a modernização de uma ferramenta interna ou a necessidade de uma solução rápida mas confiável para o processamento em background, o Postgres pode ser o herói da sua job queue. Com um schema bem projetado, a mágica do FOR UPDATE SKIP LOCKED e um pouco de lógica nos workers, você terá um sistema que mantém as tarefas organizadas, os workers ocupados e a infraestrutura simplificada. Aproveite a simplicidade e a confiabilidade do seu novo sistema de filas, movido pelo bom e velho Postgres!


Quer mais?

Entre na PATREON Community para apoiar o projeto e ter acesso a conteúdos premium como artigos, vídeos e insights variados. Lembre-se de entrar na comunidade do PATREON para receber informações valiosas, tutoriais, insights, suporte prioritário e muito mais. Você também pode apoiar o projeto pelo Buy Me a Coffe e ter os mesmos benefícios.

Os assinantes do PATREON têm acesso a um projeto Delphi completo que implementa um sistema de job queue.

No artigo do PATREON você encontra 2 projetos:

  • JobProducer.dproj: é o programa cliente que pede a execução de algum job. No exemplo ele pode pedir 2 tipos de job (“send_email”, “create_report”), mas você pode estendê-lo para o que precisar.
  • QueueWorker.dproj: é o worker propriamente dito. Você pode abrir várias instâncias e verificar que o sistema atribui cada job a uma única instância de worker. Se o worker travar depois de pegar um job, esse job volta a ficar disponível para os outros workers.

Bom proveito!

– Daniele Teti

Comments