CodeGym /Cursos /JAVA 25 SELF /Pipelines de processamento de arquivos: producer–consumer...

Pipelines de processamento de arquivos: producer–consumer

JAVA 25 SELF
Nível 59 , Lição 4
Disponível

1. Como coordenar várias threads no processamento de arquivos

Quando você trabalha com arquivos grandes ou com processamento de dados complexo, é comum surgir a tarefa de dividir o trabalho entre várias threads. Por exemplo, uma thread lê linhas de um arquivo e outras — processam essas linhas (procuram palavras, contam estatísticas etc.).

Qual é o problema?

  • Se todas as threads trabalharem com o mesmo recurso (por exemplo, com um arquivo ou uma coleção compartilhada), é fácil ter erros: algumas threads podem “ultrapassar” outras; surgem condições de corrida, perda de dados e sobrecarga de memória.
  • Se houver muitas threads e não houver coordenação — algumas ficarão ociosas, enquanto outras ficarão sobrecarregadas.

Precisamos de uma solução que:

  • Permita trocar dados com segurança entre threads.
  • Evite estouro de memória (se o processamento for mais lento do que a leitura).
  • Permita encerrar facilmente o trabalho de todas as threads quando os dados terminarem.

2. Padrão “Producer–Consumer” (Produtor–Consumidor)

Producer–Consumer é um padrão clássico que ajuda as threads a trabalharem em harmonia, sem atrapalhar umas às outras.

Nele há dois papéis: o Producer (produtor) cria dados — por exemplo, lê linhas de um arquivo ou recebe mensagens da rede — e os coloca em uma fila compartilhada. O Consumer (consumidor) pega os dados dessa fila e os processa: conta palavras, salva no banco de dados ou escreve em outro arquivo.

A ideia principal é que esses dois tipos de threads trabalham de forma independente. O produtor pode ler mais rápido do que o consumidor consegue processar, ou vice-versa — e, ainda assim, ninguém fica preso ao outro. Entre eles há um buffer — uma fila que equaliza o ritmo de trabalho.

Esquema visual

[Arquivo] --(lê)--> [Producer] --(coloca na fila)--> [BlockingQueue] --(pega)--> [Consumer] --(processa)

3. Implementação com BlockingQueue

No Java, o BlockingQueue (por exemplo, a implementação ArrayBlockingQueue) é perfeito para organizar essa troca.

O que é BlockingQueue?

BlockingQueue é uma fila thread-safe, com tamanho limitado, que cuida da sincronização por você. Se o produtor tentar adicionar um elemento e a fila já estiver cheia, a thread é automaticamente bloqueada e espera até surgir espaço. De forma análoga, se o consumidor tentar pegar um elemento e a fila estiver vazia, ele simplesmente espera até alguém colocar algo na fila.

Esse mecanismo resolve automaticamente o problema clássico de “corrida” entre threads e do estouro de memória: produtores não entopem a fila com elementos demais, e consumidores não tentam trabalhar com o vazio. Todas as threads trabalham de forma calma e coordenada.

Exemplo de criação da fila

import java.util.concurrent.*;

BlockingQueue<String> queue = new ArrayBlockingQueue<>(100); // buffer de 100 elementos

100 — número máximo de elementos na fila. Se o producer trabalhar mais rápido, ele vai esperar até que o consumer processe parte dos dados.

4. Coordenação de threads: backpressure e encerramento

Limitação do tamanho da fila (backpressure)

Backpressure é o mecanismo que impede o producer de “inundar” toda a memória, caso o consumer não consiga processar na mesma velocidade.

  • Se a fila estiver cheia, o producer automaticamente “freia” (o método put() bloqueia).
  • Se a fila estiver vazia, o consumer espera (o método take() bloqueia).

Isso permite que o sistema funcione de forma estável mesmo com velocidades diferentes entre as threads.

Encerramento: “poison pill” (pílula envenenada)

Quando o producer termina de ler o arquivo, é preciso avisar os consumers de que não haverá mais dados e que é hora de encerrar.

Solução:

  • Colocar na fila um objeto especial — a “poison pill” (por exemplo, a string "__END__" ou outro valor especial, mas não null).
  • Ao receber a “poison pill”, o consumer entende que é hora de encerrar.

Se houver vários consumers, coloque tantas “poison pills” quanto o número de consumers!

5. Exemplo de pipeline: leitura de arquivo e processamento de linhas

Vamos implementar um pipeline simples:

  • Uma thread lê linhas do arquivo e as coloca na fila.
  • Várias threads pegam as linhas da fila, contam a quantidade de palavras e exibem o resultado.

Passo 1. Producer — leitor de arquivo

import java.io.*;
import java.util.concurrent.*;

public class FileProducer implements Runnable {
    private final BlockingQueue<String> queue;
    private final File file;
    private final int consumerCount;
    private final String POISON_PILL = "__END__";

    public FileProducer(BlockingQueue<String> queue, File file, int consumerCount) {
        this.queue = queue;
        this.file = file;
        this.consumerCount = consumerCount;
    }

    @Override
    public void run() {
        try (BufferedReader reader = new BufferedReader(new FileReader(file))) {
            String line;
            while ((line = reader.readLine()) != null) {
                queue.put(line); // se a fila estiver cheia — aguardamos
            }
        } catch (IOException | InterruptedException e) {
            e.printStackTrace();
        } finally {
            // Colocamos uma poison pill para cada consumer
            try {
                for (int i = 0; i < consumerCount; i++) {
                    queue.put(POISON_PILL);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }
}

Passo 2. Consumer — processador de linhas

import java.util.concurrent.BlockingQueue;

public class LineConsumer implements Runnable {
    private final BlockingQueue<String> queue;
    private final String POISON_PILL = "__END__";

    public LineConsumer(BlockingQueue<String> queue) {
        this.queue = queue;
    }

    @Override
    public void run() {
        try {
            while (true) {
                String line = queue.take(); // se a fila estiver vazia — aguardamos
                if (POISON_PILL.equals(line)) {
                    break; // encerramento
                }
                int wordCount = line.trim().isEmpty() ? 0 : line.trim().split("\\s+").length;
                System.out.println(Thread.currentThread().getName() + ": " + wordCount + " palavras na linha: " + line);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

Passo 3. Inicialização do pipeline via Executors

import java.io.File;
import java.util.concurrent.*;

public class PipelineDemo {
    public static void main(String[] args) throws Exception {
        int consumerCount = 3;
        BlockingQueue<String> queue = new ArrayBlockingQueue<>(100);

        File file = new File("input.txt"); // seu arquivo

        // Iniciamos o producer
        Thread producer = new Thread(new FileProducer(queue, file, consumerCount));
        producer.start();

        // Iniciamos os consumers via ExecutorService
        ExecutorService consumers = Executors.newFixedThreadPool(consumerCount);
        for (int i = 0; i < consumerCount; i++) {
            consumers.submit(new LineConsumer(queue));
        }

        // Aguardamos o término do producer
        producer.join();

        // Aguardamos o término dos consumers
        consumers.shutdown();
        consumers.awaitTermination(1, TimeUnit.MINUTES);

        System.out.println("Processamento concluído!");
    }
}

6. Visualização: esquema de funcionamento do pipeline

flowchart LR
    A[Producer: lê o arquivo] -- put() --> Q[BlockingQueue]
    Q -- take() --> C1[Consumer 1: conta palavras]
    Q -- take() --> C2[Consumer 2: conta palavras]
    Q -- take() --> C3[Consumer 3: conta palavras]
    A -.->|poison pill| Q
    Q -.->|poison pill| C1
    Q -.->|poison pill| C2
    Q -.->|poison pill| C3

7. Dicas úteis e boas práticas

  • Limite o tamanho da fila — isso protege contra estouro de memória.
  • Use a “poison pill” para encerrar — caso contrário, os consumers podem ficar bloqueados para sempre.
  • Não use null como “poison pill” se a fila puder conter valores null reais. Prefira uma string ou um objeto especial.
  • Trate InterruptedException — isso é importante para o encerramento correto das threads.
  • Use ExecutorService para gerenciar o pool de threads — é mais simples e seguro do que criar threads manualmente.

8. Erros típicos na implementação do pipeline

Erro nº 1: fila sem limite → OutOfMemoryError. Se você usar LinkedBlockingQueue sem limite de tamanho, o producer pode “inundar” toda a memória se o consumer não acompanhar.

Erro nº 2: não colocou a “poison pill” → os consumers ficam bloqueados. Se você esquecer de colocar as “pílulas”, as threads consumidoras esperarão por novos dados para sempre.

Erro nº 3: colocou apenas uma “poison pill”, mas há vários consumers. Cada thread consumidora precisa da sua própria “pílula” — caso contrário, nem todas encerrarão.

Erro nº 4: não tratou InterruptedException. Se a thread for interrompida, é preciso encerrar corretamente (restaurar o flag de interrupção); caso contrário, podem ocorrer threads “presas”.

Erro nº 5: uso de variável compartilhada entre threads sem sincronização. Não tente “reinventar a roda” — use BlockingQueue, e não listas ou arrays comuns.

1
Pesquisa/teste
Trabalho paralelo com arquivos, nível 59, lição 4
Indisponível
Trabalho paralelo com arquivos
Trabalho paralelo com arquivos
Comentários
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION