Implementação Java com Streaming SSE para Respostas Incrementais do ChatGPT

Para reproduzir a experiência interativa do ChatGPT oficial — onde as respostas aparecem caractere por caractere em tempo real — é essencial adotar um mecanismo de transmissão contínua do servidor para o cliente. Embora a integração básica com APIs de LLMs via HTTP tradicional (como POST síncrono) seja trivial, a renderização progressiva exige uma arquitetura baseada em Server-Sent Events (SSE).

SSE é um padrão nativo do HTTP/1.1 que permite ao backend enviar múltiplos eventos textuais através de uma única conexão longa, sem necessidade de polling ou upgrade de protocolo (diferente do WebSocket). Ele opera exclusivamente sobre requisições GET, o que impõe restrições de design — especialmente quando se precisa transmitir payloads estruturados, como mensagens de usuário com contexto ou metadados.

Arquitetura adaptativa para suporte a dados de entrada

Como o EventSource do navegador não aceita corpo em requisições GET, optamos por uma estratégia de desacoplamento:

  1. O frontend envia a mensagem via POST /chat/submit, recebendo um identificador único (sessionKey) em troca.
  2. O cliente inicia então um fluxo SSE em GET /chat/stream/{sessionKey}, que aciona o processamento assíncrono da resposta.
  3. O payload original é recuperado do cache temporário (ex: ConcurrentHashMap ou Redis), e o serviço de chat é invocado em uma thread separada para evitar bloqueio no retorno do SseEmitter.

Código do controlador Spring Boot

private final Map<String, String> pendingRequests = new ConcurrentHashMap<>();
private final ChatStreamingService streamingService;

public StreamingChatController(ChatStreamingService streamingService) {
    this.streamingService = streamingService;
}

@PostMapping("/submit")
public ResponseEntity<Map<String, String>> submitQuery(@RequestBody Map<String, String> payload) {
    String query = payload.getOrDefault("prompt", "");
    String key = UUID.randomUUID().toString().replace("-", "").substring(0, 12);
    pendingRequests.put(key, query);
    return ResponseEntity.ok(Map.of("key", key));
}

@GetMapping(value = "/stream/{key}", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter startStream(@PathVariable String key) {
    String input = pendingRequests.remove(key);
    if (input == null || input.trim().isEmpty()) {
        return new SseEmitter(-1L);
    }

    SseEmitter emitter = new SseEmitter(30_000L); // timeout de 30s

    // Execução assíncrona para liberar imediatamente o endpoint
    CompletableFuture.runAsync(() -> {
        try {
            streamingService.generateResponse(input, chunk -> {
                try {
                    // Empacota cada chunk como JSON válido com campo 'text'
                    String jsonChunk = String.format("{\"text\":\"%s\"}", 
                        escapeJson(chunk.replaceAll("\n", "\\n")));
                    emitter.send(SseEmitter.event()
                        .name("delta")
                        .data(jsonChunk));
                } catch (IOException e) {
                    emitter.completeWithError(e);
                }
            });
            emitter.complete();
        } catch (Exception e) {
            emitter.completeWithError(e);
        }
    });

    return emitter;
}

private String escapeJson(String raw) {
    return raw.replace("\\", "\\\\")
              .replace("\"", "\\\"")
              .replace("\b", "\\b")
              .replace("\f", "\\f")
              .replace("\r", "\\r")
              .replace("\t", "\\t");
}

}


</div>### Cliente JavaScript com tratamento robusto de eventos

<div>```
function initiateChat(prompt) {
    fetch('/chat/submit', {
        method: 'POST',
        headers: { 'Content-Type': 'application/json' },
        body: JSON.stringify({ prompt })
    })
    .then(r => r.json())
    .then(data => {
        const streamUrl = `/chat/stream/${data.key}`;
        const es = new EventSource(streamUrl);

        es.onopen = () => console.log('Conexão SSE estabelecida');

        es.addEventListener('delta', event => {
            try {
                const parsed = JSON.parse(event.data);
                const fragment = document.createTextNode(parsed.text || '');
                document.getElementById('response-container').appendChild(fragment);
                
                // Scroll automático inteligente
                const container = document.getElementById('response-container');
                container.scrollTop = container.scrollHeight;
            } catch (e) {
                console.warn('Falha ao processar chunk:', e);
            }
        });

        es.onerror = () => {
            console.error('Erro na conexão SSE');
            es.close();
        };
    });
}
  • Thread segura: O uso de CompletableFuture.runAsync() (ou @Async) garante que o SseEmitter seja retornado imediatamente, evitando tiemouts no gateway e mantendo o fluxo ativo.
  • Escapamento JSON rigoroso: Caracteres especiais (espaços, quebras de linha, aspas) são tratados antes da serialização para preservar formatação e evitar corrupção no parser do cliente.
  • Scroll dinâmico confiável: A região de exibição deve ser um elemento <div> com overflow-y: auto, enquanto o conteúdo é injetado em um nó filho — isso resolve falhas comuns no comportamento de scrollTop em browsers modernos.
  • Timeout explícito: O construtor SseEmitter(30_000L) previne vazamentos de recursos caso o cliente feche a conexão inesperadamente.

A abordagem descrita elimina dependências externas, funciona com qualquer cliente compatível com EventSource, e oferece baixa latência visual mesmo sob carga moderada — replicando fielmente a sensação de digitação humana ao vivo.

Tags: spring-boot SSE java chatgpt-api streaming

Publicado em 7-30 18:54