rocketmq-broker: Recebe mensagens dos produtores e as armazena (através do rocketmq-store), os consumidores obtêm mensagens deste componente.
rocketmq-client: Fornece APIs de cliente para envio e recebimento de mensagens.
rocketmq-namesrv: NameServer, similar ao Zookeeper, armazena metadados em tempo de execução como TopicName e filas.
rocketmq-common: Classes, métodos e estruturas de dados genéricos.
rocketmq-remoting: Cliente/servidor baseado em Netty4 + serialização fastjson + protocolo binário personalizado.
rocketmq-store: Armazenamento de mensagens, índices, etc.
rocketmq-filtersrv: Servidor de filtragem de mensagens. Requer upload de código para o MQ para implementar filtragem.
rocketmq-tools: Ferramentas de linha de comando.
Transações Distribuídas com RocketMQ
Quando falamos em transações distribuídas, inevitavelmente encontramos o problema clássico de "transferência de conta": duas contas em bancos de dados diferentes ou subsistemas distintos, onde uma conta precisa ser debitada enquanto outra é creditada. Como garantir a atomicidade?
A abordagem comum é usar um middleware de mensagens para alcançar "consistência eventual": o Sistema A debita a conta, envia uma mensagem para o middleware, e o Sistema B recebe essa mensagem para creditar a conta.
No entanto, surge um problema: o Sistema A deve primeiro atualizar o banco de dados e depois enviar a mensagem? Ou primeiro enviar a mensagem e depois atualizar o banco de dados?
Se atualizarmos o banco de dados primeiro e o envio da mensagem falhar, mesmo com tentativas de reenvio, o que fazer? Se enviarmos a mensagem primeiro e a atualização do banco de dados falhar, a mensagem já foi enviada e não pode ser retirada. O que fazer nesse caso?
Portanto, podemos concluir: enquanto o envio de mensagens e a atualização do banco de dados não forem operações atômicas, independentemente da ordem, haverá problemas.
Solução Incorreta
Alguns podem sugerir colocar o envio de mensagens e a atualização do banco de dados na mesma transação. Se o envio falhar, a atualização do banco de dados seria automaticamente revertida. Isso garantiria a atomicidade das duas operações?
Essa abordagem parece correta, mas na verdade é errada por duas razões:
(1) O problema dos dois generais em redes: se o envio da mensagem falhar, o remetente não sabe se o middleware realmente não recebeu a mensagem ou se a recebeu, mas falhou ao enviar a resposta.
Se a mensagem foi recebida, mas o remetente acredita que não foi, e executa a reversão do banco de dados,
resultará em uma situação onde a conta A não foi debitada, mas a conta B foi creditada.
(2) Colocar chamadas de rede dentro de transações de banco de dados pode causar transações longas devido a atrasos na rede. Em casos graves, isso pode bloquear todo o banco de dados, um risco significativo.
Com base nesta análise, concluímos que esta solução está incorreta!
Solução 1 - Implementação pelo Lado do Negócio
Supondo que o middleware de mensagens não ofereça a funcionalidade de "mensagens de transação", como no caso do Kafka, como resolveríamos este problema?
A solução é a seguinte:
(1) O produtor prepara uma tabela de mensagens, colocando as operações de atualização do banco de dados e inserção de mensagem na mesma transação do banco de dados.
(2) Preparar um programa em segundo plano que continuamente transmite as mensagens da tabela local para o middleware de mensagens. Em caso de falha, tenta repetidamente. Permite mensagens duplicadas, mas não perde mensagens nem altera sua ordem.
(3) O consumidor prepara uma tabela de verificação de duplicatas. Mensagens processadas são registradas nesta tabela, implementando idempotência do negócio. No entanto, isso introduz outro problema de atomicidade: como garantir a atomicidade entre o consumo da mensagem e a inserção na tabela de verificação?
Se o consumo for bem-sucedido, mas a inserção na tabela de verificação falhar, o que fazer?
Através destes três passos, basicamente resolvemos o problema de atomicidade entre a atualização do banco de dados e o envio de mensagens de rede.
No entanto, uma desvantagem desta solução é a necessidade de projetar uma tabela de mensagens no banco de dados e um processo em segundo plano para verificar continuamente as mensagens locais. Isso acopla o processamento de mensagens com a lógica de negócio, aumentando a carga no lado do negócio.
Solução 2 - Mensagens de Transação do RocketMQ
Para resolver este problema sem acoplar com o negócio, o RocketMQ introduziu o conceito de "mensagens de transação".
p>Especificamente, o envio de mensagens é dividido em duas fases: fase de preparação e fase de confirmação. Os dois passos mencionados anteriormente são decompostos em três passos:
(1) Enviar mensagem preparada
(2) Atualizar banco de dados
(3) Com base no resultado da atualização do banco de dados (sucesso ou falha), confirmar ou cancelar a mensagem preparada.
Alguém pode perguntar: e se os dois primeiros passos forem executados com sucesso, mas o último falhar? Aqui entra um ponto chave do RocketMQ: o RocketMQ verifica periodicamente (padrão é 1 minuto) todas as mensagens preparadas, perguntando ao remetente se a mensagem deve ser confirmada ou cancelada.
A implementação do código é a seguinte:
// Listener que será chamado quando o RocketMQ encontrar mensagens preparadas
VerificadorTransacao verificadorTransacao = new VerificadorTransacaoImpl();
// Constrói o produtor de mensagens de transação
TransactionMQProducer produtor = new TransactionMQProducer("nomeGrupo");
// Define o verificador de transações
produtor.setTransactionCheckListener(verificadorTransacao);
// Lógica da transação local, equivalente à verificação da conta Bob e débito
ExecutorTransacaoImpl executorTransacao = new ExecutorTransacaoImpl();
produtor.iniciar();
// Constrói a mensagem, parâmetros omitidos
Mensagem msg = new Mensagem(......);
// Envia a mensagem
ResultadoEnvio resultadoEnvio = produtor.enviarMensagemEmTransacao(msg, executorTransacao, null);
produtor.desligar();
A execução da transação local é feita da seguinte forma:
public ResultadoEnvioTransacao enviarMensagemEmTransacao(.....) {
// 1. Envia a mensagem
ResultadoEnvio resultadoEnvio = this.enviar(msg);
// 2. Se o envio for bem-sucedido, executa a unidade de transação local
EstadoTransacaoLocal estadoTransacao = executorTransacao.executarTransacaoLocal(msg, arg);
// 3. Finaliza a transação
this.finalizarTransacao(resultadoEnvio, estadoTransacao, excecaoLocal);
}
Em resumo, comparando a Solução 2 com a Solução 1, a maior mudança do RocketMQ é transferir a responsabilidade de "verificar a tabela de mensagens" do lado do negócio para o middleware de mensagens.
Quanto à tabela de mensagens, na verdade ela ainda não é eliminada, pois o middleware precisa perguntar ao remetente se a transação foi executada com sucesso, o que requer uma "tabela de mensagens local modificada" para registrar o estado da transação.
Intervenção Manual
Alguém pode perguntar: tanto na Solução 1 quanto na Solução 2, se o remetente colocar com sucesso a mensagem na fila, mas o consumidor falhar ao processá-la, o que fazer?
Se o consumo falhar, mesmo após várias tentativas, o que fazer? Devemos reverter automaticamente todo o processo?
A resposta é intervenção manual. Do ponto de vista da prática de engenharia, o custo de reverter automaticamente todo o processo é extremamente alto, não só complexo de implementar, mas também introduz novos problemas. Por exemplo, se a reversão automática falhar, como lidar com isso?
Para casos de probabilidade extremamente baixa, a intervenção manual é mais confiável e simples do que implementar um sistema complexo de reversão automática.