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.
GO TO FULL VERSION