아키텍처 개요
가상 앵커 기반 방송 환경에서 시청자의 실시간 텍스트 입력을 즉시 처리하고 피드백하는 기능은 사용자 몰입도와 전환율에 직접적인 영향을 미칩니다. 이 시스템은 밀리초 단위로 데이터를 수집, 분석, 응답 생성, 멀티미디어 출력까지 연결하는 저지연 파이프라인으로 구성됩니다.
┌─────────────┐ ┌──────────────┐ ┌──────────────┐
│ 데이터 수집 │───▶│ 의도 분석 │───▶│ 응답 합성 │
│ (WebSocket)│ │ (패턴매칭) │ │ (템플릿/AI) │
└─────────────┘ └──────────────┘ └──────────────┘
│
┌───────────────────┘
▼
┌──────────────────┐
│ 실행 제어 및 출력 │
│ (TTS/화면연동) │
└──────────────────┘
전체 흐름은 네 개의 독립된 계층으로 분리됩니다. 각 모듈은 메시지 큐를 통해 비동기적으로 데이터를 전달하며, 비즈니스 판정 로직과 인프라 처리 로직을 명확히 격리하여 확장성을 확보합니다.
핵심 구성 요소 구현
1. 스트리밍 데이터 수집 계층
플랫폼의 WebSocket 엔드포인트를 통해 유입되는 원시 패킷을 표준화된 메시지로 변환하는 역할만 수행합니다. 외부 서비스 의존성을 최소화하고 메모리 부하를 줄이기 위해 파싱 외의 조건부 판정 로직은 포함하지 않습니다.
import asyncio
from typing import AsyncIterator, Callable
class LiveCommentIngestor:
def __init__(self, ws_endpoint: str, filters: list[Callable]):
self.endpoint = ws_endpoint
self.filters = filters
self._downstream_handlers = []
def attach_processor(self, callback: Callable):
self._downstream_handlers.append(callback)
async def fetch_stream(self) -> AsyncIterator[dict]:
# 실제 WebSocket 연결 및 패킷 수신 구현부
pass
async def execute(self):
async for raw_payload in self.fetch_stream():
structured_data = self._transform_packet(raw_payload)
if not structured_data or structured_data.get('msg_type') != 'text':
continue
for validator in self.filters:
if not validator(structured_data):
break
else:
for handler in self._downstream_handlers:
await handler(structured_data)
def _transform_packet(self, data: bytes) -> dict:
return {
'viewer_id': data.get('uid'),
'message': data.get('txt'),
'event_time': data.get('ts')
}
2. 의도 분석 라우터
수신된 텍스트를 사전 정의된 비즈니스 카테고리(가격 문의, 사이즈 상담, 배송 추적, 품질 검증, 일반 대화)로 매핑합니다. 운영 담당자가 설정 파일을 기반으로 키워드 사전을 실시간으로 갱신할 수 있도록 컴파일된 정규식 패턴을 캐싱하여 처리 효율을 높입니다.
import re
from dataclasses import dataclass
@dataclass
class RuleConfig:
category_label: str
search_patterns: list[str]
class IntentRouter:
def __init__(self, rule_configs: list[RuleConfig]):
self.compiled_rules = [
(cfg.category_label, [re.compile(p) for p in cfg.search_patterns])
for cfg in rule_configs
]
self.default_label = "general_chat"
def evaluate(self, input_text: str) -> str:
normalized = input_text.lower().strip()
for category, patterns in self.compiled_rules:
if any(p.search(normalized) for p in patterns):
return category
return self.default_label
3. 응답 합성 엔진
분석된 의도 태그에 따라 적절한 텍스트를 산출합니다. 빈도가 높고 정형화된 질문은 로컬 캐시된 템플릿으로 즉시 처리하며, 복합적이거나 문맥이 필요한 질문은 외부 LLM API로 라우팅합니다. 이 전략은 처리 지연 시간을 단축하고 외부 호출 비용을 통제합니다.
class ResponseComposer:
def __init__(self, template_store: dict, ai_proxy):
self.static_cache = template_store
self.external_model = ai_proxy
async def assemble_output(self, intent_tag: str, metadata: dict) -> str:
if intent_tag in self.static_cache:
tpl_string = self.static_cache[intent_tag]
return tpl_string.format_map(metadata)
try:
return await self.external_model.generate(
query=metadata['message'],
context={"intent": intent_tag}
)
except Exception:
return "시스템 점검 중입니다. 잠시 후 다시 시도해 주세요."
4. 실행 제어 계층
생성된 텍스트 답변을 가상 인간의 음성 출력(TTS), 입모양 동기화(Lip-sync), 보조 화면(KT판/이미지) 표시 명령으로 변환합니다. 각 출력 채널은 별도의 스레드 풀에서 병렬로 처리되며, 렌더링 완료 신호를 대기열에 반환합니다.
지연 시간 최적화 기법
실시간 상호작용의 핵심 성능 지표는 E2E(End-to-End) 지연 시간입니다. 댓글 입력부터 가상 인간이 음성을 시작하기까지의 시간은 3초 이내로 관리되어야 하며, 이를 위해 다음과 같은 기술적 조치가 필요합니다.
- 스트리밍 TTS 적용: 전체 문장 생성을 기다리지 않고 청크 단위로 오디오 버퍼를 재생하여 첫 바이트 응답 시간(First Byte Latency)을 단축합니다.
- 프리페치 캐싱: 자주 호출되는 템플릿 답변에 대응하는 오디오 파일을 사전에 렌더링하여 메모리에 상주시킵니다.
- 큐 우선순위 분할: 메시지 큐를 VIP(구매 의도), NORMAL(문의), LOW(스팸/잡담)으로 분할하여 리소스 독점을 방지하고 고가치 응답을 우선 스케줄링합니다.