O ApsaraMQ for RocketMQ entrega mensagens pelo menos uma vez. Isso significa que um consumidor pode receber a mesma mensagem mais de uma vez. Se a lógica do seu negócio for sensível a duplicatas — como débitos de pagamento, ajustes de estoque ou criação de pedidos — implemente o consumo idempotente. O consumo idempotente garante que processar a mesma mensagem várias vezes produza o mesmo resultado de processá-la apenas uma vez.
Por exemplo, suponha que um consumidor processe uma mensagem de débito de pagamento para um pedido de USD 100. Devido a um problema de rede, a mensagem é entregue duas vezes. Com o consumo idempotente, o sistema debita o pagamento apenas uma vez e gera somente um registro de débito de USD 100 para o pedido.
Por que ocorrem mensagens duplicadas
Mensagens duplicadas surgem em três cenários.
Nova tentativa do produtor
Um produtor envia uma mensagem e o broker do ApsaraMQ for RocketMQ a persiste. No entanto, o broker pode falhar ao confirmar o recebimento para o produtor devido a um problema transitório de rede ou a uma falha no próprio produtor. O produtor interpreta isso como um envio com falha e tenta novamente. Consequentemente, o consumidor recebe duas mensagens com o mesmo conteúdo, mas IDs de mensagem diferentes.
Reentrega pelo broker
Um consumidor recebe e processa uma mensagem, mas a confirmação de volta para o broker falha devido a um problema transitório de rede. Como o broker não consegue confirmar se a mensagem foi consumida, ele a reentrega após a recuperação da rede para honrar a garantia de entrega pelo menos uma vez. Nesse caso, o consumidor recebe duas mensagens com o mesmo conteúdo e o mesmo ID de mensagem.
Balanceamento de carga
Eventos como instabilidades de rede, reinicializações de brokers ou reinicializações da aplicação consumidora acionam o balanceamento de carga. Durante o rebalanceamento, um consumidor pode receber mensagens já entregues anteriormente.
Use chaves de negócio, não IDs de mensagem
A primeira reação natural é deduplicar com base no ID da mensagem, mas essa abordagem não é confiável. Conforme descrito no cenário de nova tentativa do produtor, a mesma mensagem lógica pode chegar com dois IDs de mensagem diferentes quando o produtor tenta reenviá-la. A deduplicação baseada no ID da mensagem ignoraria completamente essas duplicatas.
Em vez disso, atribua um identificador único de negócio como chave da mensagem. Por exemplo, use um ID de pedido, um ID de transação de pagamento ou qualquer valor que identifique exclusivamente a operação de negócio. Essa chave permanece consistente independentemente de quantas vezes a mensagem seja enviada ou reentregue, pois deriva da sua lógica de negócio e não é gerada pelo sistema de mensagens.
Essa abordagem oferece repetibilidade previsível. Se ocorrer uma falha e a mensagem for reenviada, o mesmo contexto de negócio sempre produzirá a mesma chave, tornando a deduplicação confiável.
Implemente o consumo idempotente
Etapa 1: Defina a chave da mensagem no produtor
Anexe um identificador único de negócio como chave da mensagem ao enviá-la.
Message message = new Message();
message.setKey("ORDERID_100");
SendResult sendResult = producer.send(message);
Substitua ORDERID_100 pelo identificador único de negócio real, como um ID de pedido ou ID de transação.
Etapa 2: Recupere a chave da mensagem no consumidor
Recupere a chave da mensagem no callback do consumidor e utilize-a para impor o processamento idempotente.
consumer.subscribe("ons_test", "*", new MessageListener() {
public Action consume(Message message, ConsumeContext context) {
String key = message.getKey()
// Perform idempotent processing based on the message key that uniquely identifies your business.
}
});
Etapa 3: Imponha a idempotência com um armazenamento de deduplicação
A chave da mensagem sozinha não impede o processamento duplicado — sua aplicação deve verificar se a chave já foi processada. Um padrão comum utiliza um banco de dados relacional com uma restrição de unicidade:
Antes de processar, insira a chave da mensagem em uma tabela de deduplicação com uma restrição de unicidade na coluna da chave.
Se a inserção for bem-sucedida, processe a mensagem.
Se a inserção falhar devido a uma violação de chave primária ou restrição de unicidade, a mensagem já foi processada. Ignore-a.
Utilize o padrão inserir-e-verificar em vez de verificar-e-inserir. No padrão verificar-e-inserir, duas threads poderiam passar pela verificação antes que qualquer uma realizasse a inserção, resultando em processamento duplicado. Depender da violação de restrição de unicidade do banco de dados é uma operação atômica e segura contra condições de corrida.
Exemplo de tabela de deduplicação (MySQL):
CREATE TABLE message_dedup (
message_key VARCHAR(255) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (message_key)
);
Exemplo de lógica de consumidor idempotente:
consumer.subscribe("ons_test", "*", new MessageListener() {
public Action consume(Message message, ConsumeContext context) {
String key = message.getKey();
try {
// Attempt to insert the message key. Fails if already processed.
insertMessageKey(key);
} catch (DuplicateKeyException e) {
// Already processed. Skip.
return Action.CommitMessage;
}
// Process the business logic.
processOrder(key);
return Action.CommitMessage;
}
});
Substitua insertMessageKey e processOrder pelos seus métodos reais de banco de dados e lógica de negócio.
Para cenários de alto throughput em que um banco de dados relacional pode se tornar um gargalo, utilize uma abordagem baseada em Redis com SETNX (SET if Not eXists).
Melhores práticas
|
Prática |
Detalhes |
|
Agrupe a deduplicação e a lógica de negócio em uma única transação |
Caso a lógica de negócio falhe após a inserção do registro de deduplicação, esse registro será revertido, permitindo nova tentativa de processamento da mensagem na próxima entrega. Sem uma transação, uma operação de negócio com falha deixa o registro de deduplicação intacto, e a mensagem nunca é reprocessada. |
|
Limpe periodicamente o armazenamento de deduplicação |
A tabela de deduplicação cresce com o tempo. Defina um período de retenção (por exemplo, 7 dias) e exclua registros mais antigos em lotes para evitar crescimento ilimitado do armazenamento. |
|
Escolha o armazenamento de deduplicação adequado ao seu throughput |
Utilize um banco de dados relacional com restrições de unicidade para throughput moderado. Para cenários de alto throughput, prefira o Redis com |