Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Batch consumption

Última atualização: Jun 27, 2026

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:

  1. Busca e cache -- Threads de busca recuperam mensagens do ApsaraMQ for RocketMQ usando long polling e as armazenam em cache localmente.

  2. 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.

batch_consume

Nota

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

Termos

.

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

ConsumeMessageBatchMaxSize

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.

BatchConsumeMaxAwaitDurationInSeconds

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();
        }
    }
}
Nota

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 ConsumeMessageBatchMaxSize com um valor alto (por exemplo, 128 ou 256) e BatchConsumeMaxAwaitDurationInSeconds com 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.

Referências