O consumo em lote entrega múltiplas mensagens a uma thread consumidora em um único despacho, em vez de uma por vez. Essa abordagem reduz a sobrecarga de chamadas de procedimento remoto (RPC) para sistemas downstream e aumenta o throughput de mensagens.
Funcionamento
Um push consumer processa o consumo em lote em duas etapas:
Busca e cache -- Threads de busca recuperam mensagens do ApsaraMQ for RocketMQ usando long polling e as armazenam em cache localmente.
Despacho -- Quando as mensagens em cache atingem o limiar de tamanho do lote ou o limiar de tempo de espera (o que ocorrer primeiro), o push consumer envia o lote para uma thread consumidora processar.

O ApsaraMQ for RocketMQ suporta tanto push consumers quanto pull consumers. O consumo em lote aplica-se apenas a push consumers. Para mais informações, consulte
.
Casos de uso
O consumo em lote é mais eficaz quando os sistemas downstream se beneficiam de operações em lote. Se seu objetivo for apenas aumentar o paralelismo, considere alternativas mais simples primeiro, como adicionar instâncias consumidoras ou ajustar tamanhos de pool de threads.
Indexação em massa -- Um sistema de pedidos upstream publica mensagens de log indexadas por um cluster Elasticsearch downstream. Cada mensagem aciona uma solicitação RPC (~10 ms). Processar 10 mensagens individualmente leva 100 ms; agrupá-las em uma única chamada de indexação em massa reduz o total para ~10 ms.
Inserções em massa no banco de dados -- Uma aplicação insere registros em um banco de dados um por vez sob alta frequência de atualização, gerando carga pesada. Agrupar 10 registros por inserção e executar flush a cada 5 segundos reduz a sobrecarga de conexão e a amplificação de escrita.
Limitações
O consumo em lote tem suporte apenas via TCP. Use a edição comercial do SDK cliente TCP para Java, versão 1.8.7.3.Final ou posterior. Para notas de lançamento e instruções de download, consulte Notas de lançamento.
Tamanho máximo do lote: 1.024 mensagens.
Tempo máximo de espera entre lotes: 450 segundos.
Parâmetros
Dois parâmetros controlam o momento do despacho de um lote. O despacho ocorre quando qualquer uma das condições for atendida, o que acontecer primeiro.
|
Parâmetro |
Tipo |
Padrão |
Intervalo válido |
Descrição |
|
|
String |
32 |
1--1.024 |
Número máximo de mensagens por lote. Quando a quantidade de mensagens em cache atinge esse valor, o SDK despacha o lote imediatamente para uma thread consumidora. |
|
|
String |
0 |
0--450 |
Tempo máximo de espera em segundos. Ao término desse intervalo, o SDK despacha todas as mensagens acumuladas, mesmo que o limiar de tamanho do lote não tenha sido atingido. |
Código de exemplo
Configure o consumo em lote por meio de Properties passadas para ONSFactory.createBatchConsumer(). O callback BatchMessageListener recebe uma List<Message> contendo até ConsumeMessageBatchMaxSize mensagens.
import com.aliyun.openservices.ons.api.Action;
import com.aliyun.openservices.ons.api.ConsumeContext;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.batch.BatchConsumer;
import com.aliyun.openservices.ons.api.batch.BatchMessageListener;
import java.util.List;
import java.util.Properties;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.tcp.example.MqConfig;
public class SimpleBatchConsumer {
public static void main(String[] args) {
Properties consumerProperties = new Properties();
consumerProperties.setProperty(PropertyKeyConst.GROUP_ID, MqConfig.GROUP_ID);
consumerProperties.setProperty(PropertyKeyConst.AccessKey, MqConfig.ACCESS_KEY);
consumerProperties.setProperty(PropertyKeyConst.SecretKey, MqConfig.SECRET_KEY);
consumerProperties.setProperty(PropertyKeyConst.NAMESRV_ADDR, MqConfig.NAMESRV_ADDR);
// Set the maximum number of messages per batch.
// Default: 32. Valid values: 1 to 1024.
consumerProperties.setProperty(PropertyKeyConst.ConsumeMessageBatchMaxSize, String.valueOf(128));
// Set the maximum wait time between batches, in seconds.
// Default: 0. Valid values: 0 to 450.
consumerProperties.setProperty(PropertyKeyConst.BatchConsumeMaxAwaitDurationInSeconds, String.valueOf(10));
BatchConsumer batchConsumer = ONSFactory.createBatchConsumer(consumerProperties);
batchConsumer.subscribe(MqConfig.TOPIC, MqConfig.TAG, new BatchMessageListener() {
@Override
public Action consume(final List<Message> messages, ConsumeContext context) {
System.out.printf("Batch-size: %d\n", messages.size());
// Process messages in batches.
return Action.CommitMessage;
}
});
// Start BatchConsumer.
batchConsumer.start();
System.out.println("Consumer start success.") ;
// Wait for a fixed period to prevent the process from exiting.
try {
Thread.sleep(200000);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
Para obter o source completo, consulte a biblioteca de código no GitHub.
Para obter uma referência completa de parâmetros, consulte Métodos e parâmetros.
Melhores práticas
Ajuste o tamanho do lote e o tempo de espera em conjunto
O despacho é acionado quando ou o tamanho do lote ou o limiar de tempo de espera é atingido. Defina ambos os parâmetros para adequar-se à sua carga de trabalho:
Cenários de alto throughput -- Defina
ConsumeMessageBatchMaxSizecom um valor alto (por exemplo, 128 ou 256) eBatchConsumeMaxAwaitDurationInSecondscom um intervalo curto (por exemplo, 1--5 segundos). Isso despacha lotes frequentemente sem esperar que o lote esteja cheio.Cenários de baixo throughput -- Configure um tamanho de lote moderado (por exemplo, 32) com um tempo de espera maior (por exemplo, 10--30 segundos) para evitar o despacho de lotes muito pequenos.
Exemplo: Com ConsumeMessageBatchMaxSize definido como 128 e BatchConsumeMaxAwaitDurationInSeconds definido como 1, um lote é despachado após 1 segundo, mesmo que menos de 128 mensagens tenham sido acumuladas. Nesse caso, messages.size() no callback retorna um valor menor que 128.
Implemente idempotência de consumo
Para otimizar o consumo em lote, implemente idempotência de mensagens no cliente consumidor e garanta que cada mensagem seja processada apenas uma vez. Para mais informações, consulte Idempotência de consumo.