Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Batch consumption

Dernière mise à jour :Aug 09, 2026

La consommation par lot transmet plusieurs messages à un thread de consommateur en une seule distribution, plutôt qu'un par un. Cette approche réduit la surcharge des appels de procédure distante (RPC) vers les systèmes en aval et augmente le débit des messages.

Fonctionnement

Un consommateur push gère la consommation par lot en deux étapes :

  1. Extraction et mise en cache -- Les threads d'extraction de messages récupèrent les messages depuis ApsaraMQ for RocketMQ via un sondage long (long polling) et les mettent en cache localement.

  2. Distribution -- Lorsque les messages mis en cache atteignent le seuil de taille du lot ou le seuil de temps d'attente (selon l'événement qui survient en premier), le consommateur push soumet le lot à un thread de consommateur pour traitement.

batch_consume

Remarque

ApsaraMQ for RocketMQ prend en charge à la fois les consommateurs push et les consommateurs pull. La consommation par lot s'applique uniquement aux consommateurs push. Pour plus d'informations, consultez

Termes

.

Cas d'utilisation

La consommation par lot est particulièrement efficace lorsque les systèmes en aval tirent parti des opérations groupées. Si votre objectif est uniquement d'accroître le parallélisme, envisagez d'abord des alternatives plus simples, comme l'ajout d'instances de consommateur ou l'ajustement de la taille des pools de threads.

  • Indexation en masse -- Un système de commandes amont publie des journaux qu'un cluster Elasticsearch aval indexe. Chaque message déclenche une requête RPC (~10 ms). Le traitement individuel de 10 messages prend 100 ms ; leur regroupement en un seul appel d'indexation en masse réduit le temps total à ~10 ms.

  • Insertions de base de données en masse -- Une application insère des enregistrements dans une base de données un par un avec une fréquence de mise à jour élevée, ce qui génère une charge importante. Le regroupement de 10 enregistrements par insertion et la validation toutes les 5 secondes réduisent la surcharge de connexion et l'amplification des écritures.

Limitations

  • La consommation par lot est prise en charge uniquement via TCP. Utilisez l'édition commerciale du SDK client TCP pour Java, version 1.8.7.3.Final ou ultérieure. Pour les notes de version et les instructions de téléchargement, consultez Notes de version.

  • Taille maximale du lot : 1 024 messages.

  • Temps d'attente maximal entre les lots : 450 secondes.

Paramètres

Deux paramètres déterminent le moment où un lot est distribué. La distribution se produit dès que l'une des conditions est remplie, selon celle qui survient en premier.

Paramètre Type Valeur par défaut Plage valide Description
ConsumeMessageBatchMaxSize String 32 1–1 024 Nombre maximal de messages par lot. Lorsque le nombre de messages mis en cache atteint cette valeur, le SDK distribue immédiatement le lot à un thread de consommateur.
BatchConsumeMaxAwaitDurationInSeconds String 0 0–450 Temps d'attente maximal en secondes. Lorsque cet intervalle s'écoule, le SDK distribue tous les messages accumulés, même si le seuil de taille du lot n'a pas été atteint.

Exemple de code

Configurez la consommation par lot via Properties transmis à ONSFactory.createBatchConsumer(). Le rappel BatchMessageListener reçoit une List<Message> contenant jusqu'à ConsumeMessageBatchMaxSize messages.

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

Bonnes pratiques

Ajustez conjointement la taille du lot et le temps d'attente

Le déclenchement de la distribution intervient dès que la taille du lot ou le seuil de temps d'attente est atteint. Définissez les deux paramètres en fonction de votre charge de travail :

  • Scénarios à haut débit -- Définissez ConsumeMessageBatchMaxSize sur une valeur élevée (par exemple, 128 ou 256) et BatchConsumeMaxAwaitDurationInSeconds sur un court intervalle (par exemple, 1 à 5 secondes). Cela permet de distribuer fréquemment des lots sans attendre qu'un lot complet soit constitué.

  • Scénarios à faible débit -- Choisissez une taille de lot modérée (par exemple, 32) avec un temps d'attente plus long (par exemple, 10 à 30 secondes) afin d'éviter la distribution de lots très petits.

Exemple : Avec ConsumeMessageBatchMaxSize défini sur 128 et BatchConsumeMaxAwaitDurationInSeconds défini sur 1, un lot est distribué après 1 seconde, même si moins de 128 messages ont été accumulés. Dans ce cas, messages.size() dans le rappel renvoie une valeur inférieure à 128.

Implémentez l'idempotence de la consommation

Pour optimiser la consommation par lot, implémentez l'idempotence des messages sur votre client consommateur afin de garantir qu'un message ne soit traité qu'une seule fois. Pour plus d'informations, consultez Idempotence de la consommation.

Références