Les objets Producteur et Consommateur dans ApsaraMQ for RocketMQ sont thread-safe. Partagez une seule instance entre les threads au lieu d'en créer une par thread.
Pourquoi partager les instances
Vous pouvez déployer plusieurs instances de producteurs et de consommateurs sur un ou plusieurs brokers. Vous pouvez également exécuter plusieurs threads pour envoyer ou recevoir des messages au sein d'une instance de producteur ou de consommateur. La création d'une instance par thread gaspille des ressources et peut dégrader le débit.
Le partage d'un seul producteur ou consommateur entre plusieurs threads permet de :
Réduire la surcharge des ressources côté client et broker
Améliorer le nombre de transactions par seconde (TPS) pour l'envoi et la réception de messages
Notes d'utilisation
Ne créez pas une instance de producteur ou de consommateur pour chaque thread. Créez une seule instance et partagez-la entre les threads.
Évitez d'envoyer des messages ordonnés depuis plusieurs threads. Le broker détermine l'ordre des messages en fonction de la séquence de réception depuis un seul producteur. Lorsque plusieurs threads envoient des données simultanément, l'ordre d'arrivée au broker peut différer de l'ordre d'envoi défini dans la logique de votre application.
Producteur partagé vs producteur par thread
Correct -- créez le producteur une seule fois et partagez-le entre les threads :
Producer producer = ONSFactory.createProducer(properties);
producer.start();
// Thread 1 and Thread 2 share the same producer
Thread t1 = new Thread(() -> producer.send(msg1));
Thread t2 = new Thread(() -> producer.send(msg2));
t1.start();
t2.start();
Incorrect -- créez un producteur par thread :
// Each thread creates its own producer -- wastes resources
Thread t1 = new Thread(() -> {
Producer p = ONSFactory.createProducer(properties);
p.start();
p.send(msg);
p.shutdown();
});
Exemple de code
L'exemple suivant crée un producteur partagé et envoie des messages simultanément depuis deux threads.
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;
import java.util.Properties;
public class SharedProducer {
public static void main(String[] args) {
Properties properties = new Properties();
// The ID of the consumer group created in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.GROUP_ID, "<your-group-id>");
// Make sure the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
// and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
properties.put(PropertyKeyConst.AccessKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
properties.put(PropertyKeyConst.SecretKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
// Message send timeout in milliseconds.
properties.setProperty(PropertyKeyConst.SendMsgTimeoutMillis, "3000");
// TCP endpoint. Get this value from the TCP endpoint section on the Instance Details
// page in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.NAMESRV_ADDR, "<your-tcp-endpoint>");
final Producer producer = ONSFactory.createProducer(properties);
// Call start() once before sending any messages.
producer.start();
// Both threads share the same producer instance.
Thread thread1 = new Thread(() -> {
try {
Message msg = new Message(
"TopicTestMQ", // Topic for normal messages
"TagA", // Message tag for consumer-side filtering
"Hello MQ".getBytes() // Message body in binary format
);
SendResult result = producer.send(msg);
if (result != null) {
System.out.println("Sent successfully. Message ID: " + result.getMessageId());
}
} catch (Exception e) {
// Handle failures: retry or persist the message for later processing.
e.printStackTrace();
}
});
Thread thread2 = new Thread(() -> {
try {
Message msg = new Message("TopicTestMQ", "TagA", "Hello MQ".getBytes());
SendResult result = producer.send(msg);
if (result != null) {
System.out.println("Sent successfully. Message ID: " + result.getMessageId());
}
} catch (Exception e) {
e.printStackTrace();
}
});
thread1.start();
thread2.start();
// Call shutdown() when the producer is no longer needed to release resources.
// producer.shutdown();
}
}
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Emplacement |
|---|---|---|
<your-group-id> |
ID du groupe de consommateurs | Console ApsaraMQ for RocketMQ |
<your-tcp-endpoint> |
Endpoint TCP | Section Endpoint TCP de la page Détails de l'instance |
Points clés :
La méthode
send()utilise une transmission synchrone. Si aucune exception n'est levée, le message est envoyé avec succès.Un tag de message est un libellé (similaire à un tag Gmail) que les consommateurs utilisent pour filtrer les messages sur le broker.
Le corps du message est constitué de données binaires. ApsaraMQ for RocketMQ ne traite pas les corps de message. Le producteur et le consommateur doivent convenir des méthodes de sérialisation et de désérialisation.
Un topic utilisé pour des messages normaux ne peut pas être utilisé pour d'autres types de messages.
Mise à l'échelle au-delà d'une seule instance
Pour obtenir un débit plus élevé, déployez plusieurs instances de producteurs ou de consommateurs sur un ou plusieurs brokers. Chaque instance est partagée par un pool de threads ; ne créez pas une instance par thread, quelle que soit l'échelle.
| Échelle | Approche |
|---|---|
| Quelques threads | Créez des threads dédiés partageant un seul producteur (comme illustré ci-dessus) |
| Nombreux expéditeurs concurrents | Utilisez un pool de threads partageant un seul producteur |
| Débit très élevé | Déployez plusieurs instances de producteurs, chacune partagée par un pool de threads |