Escalabilidade do SignalR com Redis: Mecanismos de Pub/Sub para Comunicação Entre Processos

A arquitetura de escalabilidade horizontal do SignalR depende de um backplane para sincronizar o estado das conexões distribuídas em múltiplos processos ou servidores. O Redis é frequentemente utilizado como esse barramento de mensagens, aproveitando seu mecanismo nativo de Publish/Subscribe para rotear invocações em tempo real sem acoplamento direto entre os nós.

Mapeamento de Conexões e Canais Dedicados

Quando um cliente estabelece um enlace com o hub, o servidor atribui um identificador único. Para garantir que mensagens destinadas a esse cliente sejam entregues independentemente do nó que as originou, a instância que hospedou a conexão cria uma inscrição no Redis utilizando um canal nomeado com base nesse idenitficador. Dessa forma, o barramento atua como um roteador inteligente que desacopla o emissor do receptor físico.

public async Task RegistrarNovaConexaoAsync(ContextoConexao contexto)
{
    await GarantirConexaoComBarramentoAsync();

    var canalEspecifico = _fabricaDeCanais.ObterCanalPorConexao(contexto.Id);
    _repositorioLocal.Adicionar(contexto);

    // O nó atual passa a escutar mensagens direcionadas exclusivamente a este cliente
    await _assinanteRedis.SubscribeAsync(canalEspecifico, async (canal, dados) =>
    {
        var invocacao = _serializador.Desserializar<MensagemInvocacao>(dados);
        await contexto.EnviarAsync(invocacao);
    });

    if (!string.IsNullOrWhiteSpace(contexto.IdUsuario))
    {
        await InscreverCanalDeUsuarioAsync(contexto);
    }
}

Roteamento Inteligente e Publicação Condicional

O envio de mensagens segue uma lógica de verificação de localidade. Antes de publicar qualquer dado no Redis, o servidor consulta seu repositório de conexões ativas. Se o destino estiver na mesma instância, a entrega é feita diretamente pelo socket aberto, eliminando a sobrecarga de serialização e latência de rede. Caso contrário, o payload é publicado no canal correspondente, e o nó responsável pela assinatura recebe e repassa os dados ao cliente final.

public async Task RotearMensagemAsync(string idAlvo, string metodo, object[] parametros, CancellationToken token = default)
{
    if (string.IsNullOrEmpty(idAlvo))
        throw new ArgumentException("Identificador inválido", nameof(idAlvo));

    // Verifica se o cliente está conectado a esta mesma instância
    var conexaoLocal = _repositorioLocal.BuscarPorId(idAlvo);
    if (conexaoLocal != null)
    {
        // Entrega direta, evitando serialização e tráfego no barramento
        await conexaoLocal.TransmitirAsync(new PayloadInvocacao(metodo, parametros));
        return;
    }

    // Cliente remoto: publica no Redis para que o nó responsável entregue
    var payload = _serializador.Serializar(new PayloadInvocacao(metodo, parametros));
    var canalDestino = _fabricaDeCanais.ObterCanalPorConexao(idAlvo);
    await _assinanteRedis.PublishAsync(canalDestino, payload);
}

Coordanação do Ciclo de Vida e Assinaturas Globais

Além dos canais individuais, o gerenciador de ciclo de vida mantém assinaturas em tópicos globais para operações de broadcast e administração de grupos. A inicialização dessas escutas é protegida por mecanismos de concorrência para evitar inscrições duplicadas durante a partida da aplicação ou reconexões do multiplexador. A estrutura centraliza o controle de estado e garante que comandos de gerenciamento sejam propagados consistentemente.

public class CoordenadorDeBarramentoRedis<THub> : IDisposable where THub : Hub
{
    private readonly ISubscriber _barramento;
    private readonly IConnectionMultiplexer _multiplexador;
    private readonly ArmazenamentoDeConexoes _conexoesAtivas;
    private readonly ISerializer _conversor;
    private readonly SemaphoreSlim _travaInicializacao = new(1, 1);
    private bool _barramentoPronto;

    public CoordenadorDeBarramentoRedis(IConnectionMultiplexer multiplexador, ISerializer conversor)
    {
        _multiplexador = multiplexador;
        _barramento = multiplexador.GetSubscriber();
        _conversor = conversor;
        _conexoesAtivas = new ArmazenamentoDeConexoes();
    }

    private async Task InicializarEscutasGlobaisAsync()
    {
        if (_barramentoPronto) return;

        await _travaInicializacao.WaitAsync();
        try
        {
            if (_barramentoPronto) return;

            // Canal para broadcast geral (SendAll)
            await _barramento.SubscribeAsync("signalr:global", async (_, dados) =>
            {
                var msg = _conversor.Desserializar<MensagemGlobal>(dados);
                var tarefas = _conexoesAtivas
                    .ListarTodas()
                    .Where(c => !msg.IdsExcluidos.Contains(c.Id))
                    .Select(c => c.TransmitirAsync(msg.Conteudo));

                await Task.WhenAll(tarefas);
            });

            // Canal de sincronização de grupos entre nós
            await _barramento.SubscribeAsync("signalr:grupos", async (_, dados) =>
            {
                var comando = _conversor.Desserializar<ComandoGrupo>(dados);
                await ProcessarAlteracaoDeGrupoAsync(comando);
            });

            _barramentoPronto = true;
        }
        finally
        {
            _travaInicializacao.Release();
        }
    }

    private async Task ProcessarAlteracaoDeGrupoAsync(ComandoGrupo comando)
    {
        var conexao = _conexoesAtivas.BuscarPorId(comando.IdConexao);
        if (conexao == null) return; // A conexão não reside neste processo

        if (comando.Acao == AcaoGrupo.Adicionar)
            conexao.Grupos.Add(comando.NomeGrupo);
        else
            conexao.Grupos.Remove(comando.NomeGrupo);
    }

    public void Dispose()
    {
        _barramento?.UnsubscribeAll();
        _multiplexador?.Dispose();
        _travaInicializacao?.Dispose();
    }
}

A combinação de canais dedicados por identificador, verificação de localidade antes da publicação e assinaturas assíncoras gerenciadas pelo multiplexador garante que o tráfego seja distribuído eficientemente. Essa abordagem elimina a necessidade de sticky sessiosn rígidas e permite que o cluster escale horizontalmente mantendo a consistência das mensagens em tempo real.

Tags: signalr Redis pubsub aspnet-core backplane

Publicado em 8-15 01:21