Os objetos producer e consumer no ApsaraMQ for RocketMQ são thread-safe. Compartilhe uma única instância entre as threads em vez de criar uma para cada thread.
Por que compartilhar instâncias
Você pode implantar várias instâncias de producer e consumer em um ou mais brokers. Também é possível executar múltiplas threads para enviar ou receber mensagens em uma mesma instância. Criar uma instância por thread desperdiça recursos e reduz o throughput.
Compartilhar um único producer ou consumer entre várias threads:
Reduz a sobrecarga de recursos no cliente e no broker
Aumenta as transações por segundo (TPS) no envio e recebimento de mensagens
Observações de uso
Não crie uma instância de producer ou consumer para cada thread. Crie apenas uma instância e compartilhe-a entre as threads.
Evite enviar mensagens ordenadas a partir de múltiplas threads. O broker determina a ordem das mensagens com base na sequência de recebimento de um único producer. Quando várias threads enviam mensagens simultaneamente, a ordem de chegada ao broker pode diferir da ordem de envio definida na lógica da aplicação.
Producer compartilhado versus por thread
Correto -- crie o producer uma única vez e compartilhe-o entre as 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();
Incorreto -- criar um producer para cada 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();
});
Código de exemplo
O exemplo a seguir cria um producer compartilhado e envia mensagens de duas threads simultaneamente.
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();
}
}
Substitua os placeholders a seguir pelos valores reais:
|
Placeholder |
Descrição |
Onde encontrar |
|
|
ID do grupo de consumidores |
Console do ApsaraMQ for RocketMQ |
|
|
Endpoint TCP |
Seção de endpoint TCP na página Instance Details |
Pontos importantes:
O método
send()usa transmissão síncrona. Se nenhuma exceção for lançada, a mensagem foi enviada com sucesso.A tag da mensagem atua como um rótulo (semelhante às tags do Gmail) para filtragem de mensagens no broker pelos consumers.
O corpo da mensagem consiste em dados binários. Como o ApsaraMQ for RocketMQ não processa o conteúdo das mensagens, producer e consumer devem concordar sobre os métodos de serialização e desserialização.
Um topic usado para mensagens normais não aceita outros tipos de mensagem.
Dimensionamento além de uma única instância
Para obter maior throughput, implante múltiplas instâncias de producer ou consumer em um ou mais brokers. Cada instância deve ser compartilhada por um pool de threads — não crie uma instância por thread, independentemente da escala.
|
Escala |
Abordagem |
|
Poucas threads |
Crie threads dedicadas compartilhando um único producer (conforme exemplo acima) |
|
Muitos remetentes simultâneos |
Use um pool de threads compartilhando um único producer |
|
Throughput muito alto |
Implante múltiplas instâncias de producer, cada uma compartilhada por um pool de threads |