Entendendo Tecnologia Distribuída: Resolvendo Transações Distribuídas com Mensagens de Transação do RocketMQ

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.

Tags: RocketMQ transações distribuídas middleware de mensagens consistência eventual java

Publicado em 8-31 22:39