Meloming
· 75 min read

채팅 키워드알림 서비스 - 멀티테넌트 키워드 알림

조현우
조현우

CEO & Fullstack Engineer

1. 단일 테넌트로 태어난 서비스가 멀티테넌트가 될 때

라이브 방송 채팅에서 특정 단어가 등장하면 알림을 받고 싶다는 요구가 있습니다. 브랜드명이 언급됐을 때, 특정 이슈가 화제가 됐을 때, 스트리머 이름이 불릴 때입니다.

혼자 쓰는 도구라면 간단합니다. 메시지에 단어가 있는지 보고 웹훅을 쏘면 됩니다. 키워드 알림 서비스의 최초 구현이 딱 그 범위였습니다. 키워드는 KEYWORDS 환경변수, 목적지는 SLACK_WEBHOOK_URL 환경변수 하나였습니다.

다만 멀티테넌트 전환은 처음부터 다음 단계 과제로 잡아뒀습니다. 그래서 단일 테넌트 구현에서도 판정과 발송을 각각 인터페이스로 끊어두고, 컨슈머 루프를 고치지 않은 채 구현만 갈아끼울 수 있는 지점을 두 곳으로 정했습니다. 메시지를 판정하는 Matcher 와 결과를 내보내는 Notifier 입니다.

Go
// Matcher checks a message for configured keywords.
//
// Phase 2 note: multi-tenant support will require a richer return type
// that carries tenant identity alongside matched keywords. The consumer
// loop calls Match() and passes results to Notifier.Notify().
// Both interfaces are the swap points for Phase 2 without changing the loop itself.
type Matcher interface {
	Match(message string) []string
}

예고된 멀티테넌트 전환의 뼈대는 한 덩어리로 들어왔습니다. gRPC 설정 로더, MultiTenantMatcher, Discord 노티파이어, Push 노티파이어, Router, ClickHouse 라이터입니다. 조직별 설정을 밖에서 받아오고, 메시지 하나를 조직 전체에 대고 판정하고, 결과를 목적지 종류에 맞게 흩뿌리고, 남길 것을 남기는 최소 구성입니다.

그 위로 필요해진 순서대로 층이 얹혔습니다. 조직 안을 다시 그룹으로 쪼개는 그룹 단위 매칭, 백엔드 enum과 어긋난 표기를 맞추는 채널 타입 대문자 정규화, 플랫폼 공지를 걸러내는 시스템 메시지 필터, 키워드가 아니라 발화자를 지목하는 사용자 지목 알림, 같은 채널 정보가 그룹마다 중복되던 응답을 줄이는 채널 풀 정규화, 그리고 정규식 키워드와 그 라벨 축약입니다. 뒤로 갈수록 “무엇에 걸 것인가”보다 “누가 어떤 모양으로 받는가”를 다루는 층이 두꺼워집니다.

여러 조직이 각자의 키워드로 같은 스트림을 보는 순간 구조가 달라집니다.

  • 조직마다 키워드 목록이 다릅니다.
  • 한 조직 안에서도 키워드를 그룹으로 묶고 그룹마다 다른 곳으로 보냅니다.
  • 알림 대상이 Slack일 수도, Discord일 수도, 모바일 푸시일 수도 있습니다.
  • 키워드가 아니라 특정 시청자의 발화 자체를 지목할 수도 있습니다.
  • 채팅은 계속 흐르고 처리는 그 속도를 따라가야 합니다.

메시지 하나가 들어올 때마다 조직 수만큼 순회하고, 조직마다 그룹을 순회하고, 그룹마다 키워드를 순회하면 곱셈이 됩니다. 이 글은 그 곱셈을 어디서 잘랐는지 적은 기록입니다.


2. 전체 데이터 흐름

먼저 지도를 펴겠습니다.

flowchart TB
    P["채팅 서비스 워커<br/>Publisher.Publish()"] --> S["Redis Stream chat:firehose<br/>XADD MAXLEN ~ 100000"]
    S -->|"XREADGROUP COUNT 10 BLOCK 0"| C["Consumer.processMessage()"]
    B["백엔드<br/>KeywordAlertQueryService"] -->|"gRPC GetKeywordConfigs / 30s"| L["configloader.Loader<br/>인메모리 스냅샷"]
    L -->|"GetConfigs() deep copy"| C
    C --> M["MultiTenantMatcher<br/>MatchGroups / MatchAll / MatchUserGroups"]
    M --> D["dedupeGroupMatches<br/>dedupeUserGroupMatches<br/>ChannelID 기준 합집합"]
    D --> R["notifier.Router.SendAll()"]
    R --> SL["SlackNotifier<br/>Block Kit blocks"]
    R --> DC["DiscordNotifier<br/>embeds"]
    R --> PU["PushNotifier<br/>gRPC SendPushNotification"]
    M --> CH["ClickHouseWriter.Write()<br/>keyword_alert_matches"]
    C --> A["XACK<br/>defer + 독립 컨텍스트"]

입력은 Redis Streams입니다. Kafka도 Valkey Pub/Sub도 아니고 채팅 수집 서비스의 워커가 채팅 한 건마다 파이프라인으로 두 번 XADD 하는 구조입니다. 채널별 스트림 하나와 전 채널 합류 스트림 하나입니다.

Go
// 채팅 수집 서비스 워커: 스트림 퍼블리셔
const (
	firehoseKey    = "chat:firehose"
	channelMaxLen  = 1000
	firehoseMaxLen = 100000
)

키워드 알림 서비스는 이 중 chat:firehose 하나만 봅니다. 채널이 몇 개로 늘어나든 컨슈머가 붙는 스트림 키는 상수 하나입니다.

Go
// internal/consumer/consumer.go
const streamKey = "chat:firehose"

읽기 루프는 컨슈머 그룹 방식입니다.

Go
func (c *Consumer) Run(ctx context.Context) error {
	err := c.client.XGroupCreateMkStream(ctx, streamKey, c.group, "$").Err()
	if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") {
		return err
	}

	go c.logStats(ctx)

	for {
		select {
		case <-ctx.Done():
			return nil
		default:
		}

		streams, err := c.client.XReadGroup(ctx, &redis.XReadGroupArgs{
			Group:    c.group,
			Consumer: c.name,
			Streams:  []string{streamKey, ">"},
			Count:    10,
			Block:    c.blockDuration, // 0 = infinite wait in prod; short poll in tests
		}).Result()
		...
	}
}

세 가지가 눈에 띕니다.

첫째, 그룹 생성 시작점이 "$"입니다. 처음 뜰 때 밀린 백로그를 소급해서 알리지 않겠다는 뜻입니다. 몇 시간 치 지난 채팅이 한꺼번에 Slack으로 쏟아지는 사고를 여기서 막습니다.

둘째, BUSYGROUP 문자열 검사입니다. 그룹이 이미 있다는 에러는 정상 상태이므로 통과시키고 나머지 에러만 기동 실패로 처리합니다.

셋째, Block: 0은 무한 대기입니다. 프로덕션에서는 폴링이 아니라 블로킹 읽기이고 테스트에서만 withBlockDuration()을 써서 50ms 폴링으로 낮춥니다. miniredis가 컨텍스트 취소로는 블로킹을 풀어주지 않기 때문에 테스트 전용 탈출구를 남긴 것입니다.


3. 메시지 한 번에 전체 테넌트를 훑습니다

가장 단순한 멀티테넌트 구현은 조직별로 컨슈머를 두는 것입니다. 조직이 늘어날 때마다 스트림을 읽는 주체가 늘어납니다.

이 방식은 조직 수만큼 같은 메시지를 중복해서 읽습니다. 채팅 스트림은 조직과 무관하게 하나이므로 읽기는 한 번이면 충분합니다.

그래서 매처를 멀티테넌트로 만들었습니다.

Go
// MultiTenantMatcher holds a per-org KeywordMatcher and can match a message
// against all orgs in one pass.
type MultiTenantMatcher struct {
	orgMatchers map[int32]*KeywordMatcher
}

// MatchAll returns a map of orgID -> matched keywords for every org that
// has at least one keyword hit. Orgs with no matches are omitted.
func (m *MultiTenantMatcher) MatchAll(message string) map[int32][]string {
	results := make(map[int32][]string)
	for orgID, matcher := range m.orgMatchers {
		matched := matcher.Match(message)
		if len(matched) > 0 {
			results[orgID] = matched
		}
	}
	return results
}

메시지를 한 번 받아 조직 전체를 판정하고 히트가 있는 조직만 결과에 담습니다. 대부분의 채팅은 어느 키워드에도 걸리지 않으므로 결과 맵은 보통 비어 있습니다.

여기서 실제 구현의 특징 하나를 짚어야 합니다. 매처는 프로세스 수명 동안 살아 있는 객체가 아닙니다. 메시지마다 새로 만듭니다.

Go
func (c *Consumer) processMessage(ctx context.Context, msg redis.XMessage) {
	...
	configs := c.loader.GetConfigs()

	orgKeywords := make(map[int32][]string, len(configs))
	for orgID, cfg := range configs {
		if len(cfg.Keywords) > 0 && len(cfg.Groups) == 0 {
			orgKeywords[orgID] = cfg.Keywords
		}
	}
	...
	mtMatcher := matcher.NewMultiTenantMatcher(orgKeywords)

메시지 한 건을 처리할 때마다 설정 스냅샷을 읽고, 그 스냅샷으로 매처를 새로 조립합니다. 언뜻 낭비로 보이지만 이 선택이 뒤에서 두 가지 문제를 동시에 없앱니다. 하나는 설정 갱신과 매칭 사이의 경합이고(6장), 다른 하나는 정규식 재컴파일 비용입니다(4장). 차례로 이어가겠습니다.


4. keywordPattern: 세 필드가 세 가지 판정을 나눕니다

키워드 하나는 다음 구조체로 컴파일됩니다.

Go
type keywordPattern struct {
	label string
	lower string
	regex *regexp.Regexp
}

var compiledKeywordCache sync.Map

세 필드의 역할이 각각 다릅니다.

  • label: 사용자가 입력한 원본 문자열입니다. 알림 본문과 이력에 남는 이름입니다.
  • lower: 부분 문자열 비교용으로 미리 소문자화한 값입니다. 정규식 키워드에서는 비어 있습니다.
  • regex: 컴파일된 정규식입니다. 부분 문자열 키워드에서는 nil입니다.

regex != nil이냐 아니냐가 그 키워드의 성격을 정하고 나머지 한 필드는 항상 비어 있습니다.

컴파일 함수는 다음과 같습니다.

Go
func compileKeyword(keyword string) keywordPattern {
	trimmed := strings.TrimSpace(keyword)
	if cached, ok := compiledKeywordCache.Load(trimmed); ok {
		return cached.(keywordPattern)
	}

	pattern, ok := parseRegexKeyword(trimmed)
	if ok {
		if re, err := regexp.Compile(pattern); err == nil {
			compiled := keywordPattern{label: trimmed, regex: re}
			compiledKeywordCache.Store(trimmed, compiled)
			return compiled
		}
	}
	compiled := keywordPattern{label: trimmed, lower: strings.ToLower(trimmed)}
	compiledKeywordCache.Store(trimmed, compiled)
	return compiled
}

정규식으로 파싱되지 않거나 컴파일에 실패하면 조용히 부분 문자열 키워드로 강등됩니다. 잘못된 패턴 하나가 그 조직의 알림 전체를 죽이지 않게 하는 선택입니다. 다만 사용자가 오타를 냈을 때 아무 일도 일어나지 않는 게 아니라 “문자 그대로 매칭”이라는 엉뚱한 동작으로 이어집니다. 그래서 입력 검증은 프론트가 아니라 백엔드에서 따로 막습니다(6장).

표기법 파싱은 이렇게 생겼습니다.

Go
func parseRegexKeyword(keyword string) (string, bool) {
	if strings.HasPrefix(keyword, "regex:") {
		pattern := strings.TrimPrefix(keyword, "regex:")
		return pattern, pattern != ""
	}

	if len(keyword) < 2 || keyword[0] != '/' {
		return "", false
	}
	closing := strings.LastIndex(keyword[1:], "/")
	if closing < 0 {
		return "", false
	}
	closing += 1
	pattern := keyword[1:closing]
	flags := keyword[closing+1:]
	if pattern == "" || !isSupportedSlashFlags(flags) {
		return "", false
	}
	if flags == "" {
		return pattern, true
	}
	return "(?" + uniqueFlags(flags) + ")" + pattern, true
}

두 가지 표기를 지원합니다. regex: 접두사 형태와 슬래시로 감싸는 형태입니다. 표기 자체가 의도를 담고 있어서 설정 화면에 “이건 정규식입니다” 체크박스를 따로 둘 필요가 없습니다.

슬래시 형태에서 눈여겨볼 점은 닫는 슬래시를 LastIndex로 찾는다는 것입니다. 패턴 안에 슬래시가 들어가도 마지막 슬래시를 구분자로 삼습니다. 그 뒤에 남는 문자열이 플래그이고 허용 집합은 i, m, s 셋뿐입니다.

Go
func isSupportedSlashFlags(flags string) bool {
	for _, flag := range flags {
		if flag != 'i' && flag != 'm' && flag != 's' {
			return false
		}
	}
	return true
}

Go의 regexp에는 자바스크립트처럼 패턴과 플래그를 따로 받는 API가 없습니다. 그래서 플래그를 RE2의 인라인 플래그 문법으로 바꿔 패턴 앞에 붙입니다. /hello/i(?i)hello가 됩니다. uniqueFlags가 중복 문자를 제거하는 이유는 (?ii) 같은 표현이 컴파일 에러를 내기 때문입니다.

그리고 캐시입니다.

Go
var compiledKeywordCache sync.Map

정규식 컴파일은 비쌉니다. 이 서비스에서는 그 비싼 일이 자주 일어납니다. 3장에서 본 것처럼 매처가 메시지마다 새로 만들어지고, 그룹 매칭 경로는 아예 루프 안에서 키워드마다 compileKeyword를 호출합니다.

Go
for _, kw := range g.Keywords {
	pattern := compileKeyword(kw)
	...
}

캐시가 없으면 채팅 한 줄마다 모든 조직의 모든 그룹의 모든 정규식을 다시 컴파일하게 됩니다. sync.Map을 쓴 이유는 쓰기보다 읽기가 훨씬 많고 키 집합이 설정 개수만큼으로 사실상 고정되기 때문입니다. 키는 트림된 키워드 문자열이므로 조직이 달라도 같은 키워드는 같은 컴파일 결과를 공유합니다.

정규식에는 원본을, 부분 문자열에는 소문자를

키워드 매칭은 대소문자를 구분하지 않아야 합니다. 순진하게 짜면 키워드마다 양쪽을 소문자로 바꿉니다. 키워드 개수만큼 같은 메시지를 반복해서 변환하게 됩니다.

그래서 소문자 변환의 위치를 옮겼습니다. 키워드는 컴파일 시점에 한 번, 메시지는 매칭 시작 시점에 한 번입니다.

Go
func (m *KeywordMatcher) Match(message string) []string {
	lower := strings.ToLower(message)
	var matched []string
	for _, p := range m.patterns {
		if p.regex != nil {
			if p.regex.MatchString(message) {
				matched = append(matched, p.label)
			}
			continue
		}
		if strings.Contains(lower, p.lower) {
			matched = append(matched, p.label)
		}
	}
	return matched
}

여기가 이 서비스에서 가장 조용하고 가장 중요한 두 줄입니다.

Go
if p.regex.MatchString(message)      // 원본
if strings.Contains(lower, p.lower)  // 소문자

같은 함수 안에서 같은 메시지를 두 가지 형태로 씁니다. 정규식에는 message를, 부분 문자열 비교에는 lower를 넘깁니다.

이유는 대소문자 구분의 주인이 다르기 때문입니다. 부분 문자열 키워드는 “대소문자 구분 없음”이 서비스가 정한 규칙입니다. 반면 정규식 키워드는 사용자가 /patterns/i처럼 플래그로 직접 지정합니다. 미리 소문자로 만든 문자열을 정규식에 넘기면 플래그를 붙이지 않은 사용자가 의도한 대소문자 구분이 조용히 사라집니다. /[A-Z]{3}/처럼 대문자를 잡으려는 패턴은 소문자화된 입력에서 영원히 매칭되지 않습니다.

성능을 위해 만든 값을 아무 데나 재사용하면 의미가 바뀝니다. 최적화가 닿아도 되는 범위를 코드에서 분명히 끊어 둡니다.

같은 규칙이 그룹 매칭 경로에도 그대로 복제되어 있습니다.

Go
func (m *MultiTenantMatcher) MatchGroups(
	orgID int32,
	groups []model.GroupConfig,
	message string,
) []GroupMatch {
	var out []GroupMatch
	lower := strings.ToLower(message)
	for _, g := range groups {
		var matched []string
		for _, kw := range g.Keywords {
			pattern := compileKeyword(kw)
			if pattern.regex != nil {
				if pattern.regex.MatchString(message) {
					matched = append(matched, pattern.label)
				}
				continue
			}
			if strings.Contains(lower, pattern.lower) {
				matched = append(matched, pattern.label)
			}
		}
		...
	}
}

lower는 그룹 루프 바깥에서 한 번만 만듭니다. 조직이 열 개이고 그룹이 조직마다 다섯 개여도 소문자 변환은 메시지당 두 번(경로별로 한 번씩)입니다.


5. 매칭 진입점이 네 개인 이유

matcher 패키지의 공개 메서드는 네 개입니다. 이름이 비슷해서 헷갈리기 쉬우므로 정리합니다.

메서드수신자입력반환쓰이는 곳
Match*KeywordMatcher메시지[]string 매칭 키워드단일 키워드 목록 판정
MatchAll*MultiTenantMatcher메시지map[int32][]string레거시 경로(그룹 없는 조직)
MatchGroups*MultiTenantMatcher그룹 목록 + 메시지[]GroupMatch기본 경로(그룹 있는 조직)
MatchUserGroups*MultiTenantMatcher사용자 그룹 목록 + 플랫폼 + userID[]UserGroupMatch사용자 지목 알림

MatchGroupsMatchUserGroupsMultiTenantMatcher의 메서드이면서 orgMatchers 필드를 전혀 쓰지 않는다는 점은 솔직히 말해 구조가 자란 흔적입니다. 그룹 목록을 인자로 받아 그 자리에서 컴파일하기 때문에 수신자는 사실상 네임스페이스 역할만 합니다. 테스트가 이 사실을 그대로 드러냅니다.

Go
func newTestMatcher(t *testing.T) *MultiTenantMatcher {
	t.Helper()
	return NewMultiTenantMatcher(map[int32][]string{})
}

빈 맵으로 만든 매처로 그룹 매칭을 테스트합니다. 리팩터링 대상이라는 신호이지만 동작은 정확하고 경계는 명확하므로 그대로 두었습니다.

사용자 지목 경로는 판정 규칙이 다릅니다.

Go
// Platform comparison is case-insensitive. UserID is exact match (case-sensitive)
// because platform user IDs are opaque tokens.
func (m *MultiTenantMatcher) MatchUserGroups(
	orgID int32,
	groups []model.UserGroupConfig,
	platform string,
	userID string,
) []UserGroupMatch {
	if userID == "" {
		return nil
	}
	var out []UserGroupMatch
	for _, g := range groups {
		var matched []model.AlertUser
		for _, u := range g.Users {
			if strings.EqualFold(u.Platform, platform) && u.PlatformUserID == userID {
				matched = append(matched, u)
			}
		}
		...
	}
}

플랫폼 이름은 strings.EqualFold로 대소문자를 무시하고, 플랫폼 사용자 ID는 정확히 일치해야 합니다. 플랫폼 ID는 사람이 읽으라고 만든 문자열이 아니라 불투명 토큰입니다. 여기서 정규화를 하면 엉뚱한 사람을 지목하게 될 수 있습니다. userID == ""를 먼저 걸러내는 것도 같은 이유입니다. 빈 문자열이 등록 데이터의 빈 값과 우연히 맞아떨어지는 사고를 막습니다.


6. 설정의 모양: 조직, 그룹, 채널

멀티테넌트 설정은 3층입니다. Go 쪽 모델은 다음과 같습니다.

Go
// internal/model/config.go

// ChannelConfig describes a single alert destination with full metadata.
type ChannelConfig struct {
	ChannelID  int32
	Name       string
	Type       string // "SLACK" | "DISCORD" | "PUSH"
	WebhookURL string
}

// GroupConfig holds the keyword-alert configuration for a single group
// within an organisation.
type GroupConfig struct {
	GroupID  int32
	Name     string
	IsActive bool
	Keywords []string
	Channels []ChannelConfig
}

// AlertUser identifies a specific viewer inside a UserGroupConfig.
type AlertUser struct {
	Platform       string
	PlatformUserID string
	DisplayLabel   string
}

type UserGroupConfig struct {
	GroupID  int32
	Name     string
	IsActive bool
	Users    []AlertUser
	Channels []ChannelConfig
}

조직 단위 묶음은 로더 쪽에 있습니다.

Go
// internal/configloader/loader.go

// AlertChannelConfig describes a single alert destination.
// Kept for backward compatibility with notifier.Router.SendAll.
type AlertChannelConfig struct {
	Type       string
	WebhookURL string
}

type OrgConfig struct {
	OrganizationID int32
	IsActive       bool
	Groups         []model.GroupConfig     // primary path (new): keyword-target alerts
	UserGroups     []model.UserGroupConfig // user-target alerts

	// legacy: populated by loader as fallback when Groups is empty
	Keywords []string
	Channels []AlertChannelConfig
}

ChannelConfigAlertChannelConfig가 둘 다 존재하는 이유는 이렇습니다. 전자는 ChannelID를 들고 있어서 그룹 간 중복 제거의 키가 되고, 후자는 Router가 받는 최소 형태입니다. 컨슈머가 둘 사이를 변환합니다.

Go
func modelChannelToAlert(ch model.ChannelConfig) configloader.AlertChannelConfig {
	return configloader.AlertChannelConfig{
		Type:       ch.Type,
		WebhookURL: ch.WebhookURL,
	}
}

원본 스키마는 백엔드의 Prisma 모델입니다. 조직 하나에 설정 하나, 설정 하나에 채널 여러 개, 그룹과 채널은 다대다 조인 테이블로 연결됩니다.

Prisma
model KeywordAlertConfig {
  id             Int      @id @default(autoincrement())
  organizationId Int      @unique
  isActive       Boolean  @default(true)
  // legacy: kept nullable for safe rollback. dropped in cleanup PR.
  keywords       Json?

  organization Organization            @relation(...)
  channels     KeywordAlertChannel[]
  groups       KeywordAlertGroup[]
  userGroups   KeywordAlertUserGroup[]
}

model KeywordAlertChannel {
  id         Int              @id @default(autoincrement())
  configId   Int
  name       String?          @db.VarChar(100)
  type       AlertChannelType
  webhookUrl String?
  isActive   Boolean          @default(true)
}

model KeywordAlertGroup {
  id       Int     @id @default(autoincrement())
  configId Int
  name     String  @db.VarChar(100)
  keywords Json    @default("[]")
  isActive Boolean @default(true)

  channelLinks KeywordAlertGroupChannel[]
}

model KeywordAlertGroupChannel {
  id        Int @id @default(autoincrement())
  groupId   Int
  channelId Int

  @@unique([groupId, channelId])
  @@index([channelId])
}

model KeywordAlertUserGroupUser {
  id             Int    @id @default(autoincrement())
  groupId        Int
  platform       String @db.VarChar(50)
  platformUserId String @db.VarChar(200)
  displayLabel   String? @db.VarChar(200)

  @@unique([groupId, platform, platformUserId])
}

그룹과 채널이 조인 테이블로 연결되어 있으므로 한 채널을 여러 그룹이 공유할 수 있습니다. 그래서 7장의 중복 제거가 필요해집니다.

gRPC 응답은 처음에 그룹마다 채널 객체를 통째로 실어 보냈습니다. 채널 하나를 다섯 그룹이 참조하면 같은 웹훅 URL이 다섯 번 실립니다. 그래서 채널 풀 방식으로 정규화했습니다. 조직당 채널 목록은 한 번만 싣고 그룹은 ID로 참조합니다.

TypeScript
// 백엔드: findAllActiveConfigs()
return {
  organization_id: config.organizationId,
  is_active: config.isActive,
  keywords: legacyKeywords,
  channels: channelPool,     // legacy top-level mirror
  channel_pool: channelPool,
  groups,                    // groups[].channel_ids
  user_groups: userGroups,   // user_groups[].channel_ids
};

클라이언트는 새 필드와 옛 필드를 동시에 읽을 수 있어야 배포 순서에 자유로워집니다. 그래서 로더에 dual-read 함수를 뒀습니다.

Go
// resolveChannels picks the new (channel_ids + pool) path when populated,
// falling back to the deprecated embedded channels for transitional servers.
resolveChannels := func(ids []int32, embedded []*businessv1.AlertChannel) []model.ChannelConfig {
	if len(ids) > 0 {
		out := make([]model.ChannelConfig, 0, len(ids))
		for _, id := range ids {
			if ch, ok := channelPool[id]; ok {
				out = append(out, ch)
			}
		}
		return out
	}
	out := make([]model.ChannelConfig, 0, len(embedded))
	for _, ch := range embedded {
		out = append(out, model.ChannelConfig{
			ChannelID:  ch.GetChannelId(),
			Name:       ch.GetName(),
			Type:       ch.GetType(),
			WebhookURL: ch.GetWebhookUrl(),
		})
	}
	return out
}

channelPool에 없는 ID는 조용히 버립니다. 백엔드가 비활성 채널을 풀에서 빼기 때문에 그 채널을 참조하는 링크는 자연스럽게 사라집니다. 필터링 책임을 한쪽에 몰아두고 다른 쪽은 참조 해석만 합니다.

정규식 검증도 백엔드가 맡습니다. 여기서 재미있는 크로스 런타임 제약이 나옵니다. 설정 화면은 TypeScript이고 매칭은 Go인데, Go의 RE2는 lookaround와 backreference를 지원하지 않습니다. JS RegExp로만 검증하면 저장은 되는데 매칭은 안 되는 키워드가 생깁니다.

TypeScript
function hasUnsupportedGoRegexSyntax(pattern: string): boolean {
  return (
    pattern.includes('(?=') ||
    pattern.includes('(?!') ||
    pattern.includes('(?<=') ||
    pattern.includes('(?<!') ||
    /(^|[^\\])\\[1-9]/.test(pattern)
  );
}

저장 시점에 거절하고, 메시지로 이유를 알려줍니다. 실행 런타임이 다른 두 시스템 사이에서는 “문법적으로 맞는가”가 아니라 “저쪽에서도 도는가”를 물어야 합니다.

설정은 30초마다 통째로 갈아끼웁니다

설정 로더는 gRPC 폴링입니다. 스트리밍이나 변경 알림 방식이 아니라 단순한 주기 조회입니다.

Go
func (l *Loader) Run(ctx context.Context) {
	ticker := time.NewTicker(l.pollInterval)
	defer ticker.Stop()

	for {
		select {
		case <-ctx.Done():
			return
		case <-ticker.C:
			l.fetch(ctx)
		}
	}
}

주기는 CONFIG_POLL_INTERVAL 환경변수이고 기본값은 30초입니다. 파싱에 실패하면 30초로 되돌립니다. 즉 설정 화면에서 키워드를 바꾸면 최대 30초 뒤부터 알림에 반영됩니다.

호출부에는 별도 데드라인이 있습니다.

Go
// Per-call deadline prevents polling stalls when the backend is slow/hung.
callCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()

client := businessv1.NewKeywordAlertQueryServiceClient(conn)
resp, err := client.GetKeywordConfigs(callCtx, &businessv1.GetKeywordConfigsRequest{})

폴링 주기가 30초인데 호출이 60초 걸리면 폴링이 밀립니다. 주기적 작업에는 주기보다 짧은 데드라인이 필요합니다.

갱신은 맵 통째 교체입니다.

Go
func (l *Loader) updateConfigs(configs []*OrgConfig) {
	m := make(map[int32]*OrgConfig, len(configs))
	for _, c := range configs {
		m[c.OrganizationID] = c
	}

	l.mu.Lock()
	l.configs = m
	l.mu.Unlock()
}

부분 갱신이 아니라 새 맵을 만들어 포인터만 바꿉니다. 잠금 구간이 대입 한 줄이고 삭제된 조직이 남는 문제도 자동으로 사라집니다.

읽기는 깊은 복사입니다.

Go
// GetConfigs returns a deep-copied snapshot of the current configs.
// Safe for concurrent use.
func (l *Loader) GetConfigs() map[int32]*OrgConfig {
	l.mu.RLock()
	defer l.mu.RUnlock()

	snapshot := make(map[int32]*OrgConfig, len(l.configs))
	for id, cfg := range l.configs {
		cp := &OrgConfig{...}
		copy(cp.Keywords, cfg.Keywords)
		copy(cp.Channels, cfg.Channels)
		// Groups, UserGroups도 슬라이스까지 복사
		...
	}
	return snapshot
}

여기서 3장의 “메시지마다 매처를 다시 만든다”가 의미를 얻습니다. 갱신 중 매칭 유실을 막는 장치는 따로 없습니다. 대신 구조 자체가 그 문제를 만들지 않습니다.

  • 컨슈머는 메시지 처리 시작 시점에 스냅샷을 한 번 받습니다.
  • 그 스냅샷은 깊은 복사본이므로 처리 도중 폴링이 맵을 교체해도 이미 손에 든 데이터는 바뀌지 않습니다.
  • 처리 중인 메시지는 옛 설정으로, 다음 메시지는 새 설정으로 판정됩니다. 한 메시지가 반쯤 옛 설정, 반쯤 새 설정으로 처리되는 상태는 존재하지 않습니다.

교체 가능한 상태를 오래 들고 있지 않으면 교체 시점을 걱정할 일도 줄어듭니다.

기동 시퀀스에는 반대 방향의 판단이 들어갔습니다. 설정을 못 받으면 아예 뜨지 않습니다.

Go
// Require first successful config fetch before consuming. This prevents
// silent message loss when the backend gRPC is unreachable at startup.
// Retry every 2s up to 30s total before failing.
initCtx, initCancel := context.WithTimeout(ctx, 30*time.Second)
var initErr error
for {
	initErr = loader.FetchOnceWithError(initCtx)
	if initErr == nil {
		break
	}
	slog.Warn("initial config fetch failed, retrying", "err", initErr)
	select {
	case <-initCtx.Done():
		initCancel()
		log.Fatalf("config loader startup failed after 30s: %v", initErr)
	case <-time.After(2 * time.Second):
	}
}

빈 설정으로 기동하면 서비스는 “정상”입니다. 헬스체크도 통과하고 메시지도 잘 소비하고 ACK도 잘 합니다. 그저 아무도 알림을 받지 못할 뿐입니다. 이런 종류의 정상은 장애보다 나쁩니다. 그래서 fail-closed로 바꿨습니다. 같은 판단을 ClickHouse 연결에도 적용했습니다.

Go
// ClickHouse (opt-in via ADDR). When configured, fail startup on connection
// failure to avoid silently losing history for the whole process lifetime.
if cfg.ClickHouseAddr != "" {
	chWriter, err = store.NewClickHouseWriter(...)
	if err != nil {
		log.Fatalf("clickhouse connection failed (fail-closed): %v", err)
	}
}

주소를 설정하지 않으면 이력 적재를 아예 하지 않습니다(opt-in). 하지만 주소를 설정했는데 연결이 안 되면 기동을 멈춥니다. “쓰기로 했으면 반드시 쓴다”와 “안 쓰기로 하면 안 쓴다” 사이에 애매한 상태를 두지 않습니다.


7. 중복 제거는 매처가 아니라 호출자의 몫

MatchGroups 는 활성 그룹을 돌면서 키워드가 하나라도 걸린 그룹만 돌려줍니다. 그룹 사이의 중복 제거는 하지 않습니다. 그 책임은 호출자에게 남겨뒀습니다.

같은 메시지가 두 그룹에 동시에 걸릴 수 있고, 두 그룹이 같은 채널을 가리킬 수도 있습니다. 조인 테이블 구조상 이건 예외가 아니라 흔한 경우입니다. 마케팅 그룹과 CS 그룹이 같은 Slack 채널을 보고 있으면 한 메시지에 웹훅이 두 번 날아갑니다.

이때 알림을 한 번만 보낼지 두 번 보낼지는 매칭의 문제가 아니라 전달 정책의 문제입니다. 매처는 무엇이 걸렸는지만 정확히 답하고 어떻게 보낼 것인가는 위층에서 정합니다.

호출자 쪽 구현은 다음과 같습니다.

Go
// dedupeGroupMatches combines per-group channels into a single unique slice
// (keyed by ChannelID) and unions the matched keywords across groups.
func dedupeGroupMatches(matches []matcher.GroupMatch) (channels []model.ChannelConfig, keywords []string) {
	seen := make(map[int32]struct{})
	keywordSet := make(map[string]struct{})
	for _, gm := range matches {
		for _, ch := range gm.Channels {
			if _, ok := seen[ch.ChannelID]; ok {
				continue
			}
			seen[ch.ChannelID] = struct{}{}
			channels = append(channels, ch)
		}
		for _, kw := range gm.Keywords {
			keywordSet[kw] = struct{}{}
		}
	}
	for kw := range keywordSet {
		keywords = append(keywords, kw)
	}
	return
}

채널은 ChannelID 기준 합집합, 키워드도 합집합입니다. 결과는 “채널당 알림 한 번, 본문에는 걸린 키워드 전부”입니다. 두 그룹이 각각 다른 키워드로 같은 채널을 가리켰다면 그 채널은 두 키워드가 함께 적힌 알림 하나를 받습니다.

사용자 지목 경로도 같은 형태이지만 라벨의 출처가 다릅니다.

Go
func dedupeUserGroupMatches(matches []matcher.UserGroupMatch) (channels []model.ChannelConfig, labels []string) {
	seen := make(map[int32]struct{})
	labelSet := make(map[string]struct{})
	for _, ugm := range matches {
		for _, ch := range ugm.Channels {
			...
		}
		labelSet[ugm.Name] = struct{}{}
	}
	...
}

키워드 경로의 라벨은 매칭된 키워드이고, 사용자 경로의 라벨은 그룹 이름입니다. “무엇이 걸렸는가”의 답이 경로마다 다르기 때문입니다.

두 경로는 끝까지 합쳐지지 않습니다.

Go
// Per-org matching results: keyword and user paths tracked independently
// because alert content (matched keywords vs. user-group labels) differs.
type orgResult struct {
	labels   []string
	channels []configloader.AlertChannelConfig
}
orgKeywordResults := make(map[int32]orgResult, len(configs))
orgUserResults := make(map[int32]orgResult, len(configs))

한 메시지가 키워드에도 걸리고 사용자 지목에도 걸리면 알림은 두 개 갑니다. 같은 채널이라도 그렇습니다. 걸린 이유가 다르면 받는 사람도 두 줄로 보는 편이 낫다는 판단입니다. 중복 제거는 “같은 이유” 안에서만 합니다.


8. Router.SendAll: 한 채널의 실패가 나머지를 막지 않게

알림 대상은 종류가 다릅니다. Slack과 Discord는 웹훅 URL로 보내고, 푸시는 조직 단위로 대상을 찾아 보냅니다. 인터페이스가 다르므로 둘로 나눴습니다.

Go
// Sender sends a notification to a webhook-based destination.
type Sender interface {
	Send(webhookURL string, msg model.ChatMessage, keywords []string) error
}

// PushSender sends push notifications for a given organization.
type PushSender interface {
	SendPush(ctx context.Context, orgID int32, msg model.ChatMessage, keywords []string) error
}

라우터 본체는 짧습니다.

Go
// SendResult tracks success and failure counts across all channels.
type SendResult struct {
	Success int
	Failure int
}

// SendAll iterates the alert channels and dispatches each to the correct
// notifier. Errors are logged but not returned: a failure on one channel
// must not block delivery to the others. Returns counts for metrics.
func (r *Router) SendAll(ctx context.Context, orgID int32, channels []configloader.AlertChannelConfig, msg model.ChatMessage, keywords []string) SendResult {
	result := SendResult{}
	for _, ch := range channels {
		var err error
		switch strings.ToUpper(ch.Type) {
		case "SLACK":
			err = r.slack.Send(ch.WebhookURL, msg, keywords)
		case "DISCORD":
			err = r.discord.Send(ch.WebhookURL, msg, keywords)
		case "PUSH":
			err = r.push.SendPush(ctx, orgID, msg, keywords)
		default:
			log.Printf("[router] unknown channel type %q (org=%d)", ch.Type, orgID)
			continue
		}
		if err != nil {
			log.Printf("[router] %s send error (org=%d): %v", ch.Type, orgID, err)
			result.Failure++
		} else {
			result.Success++
		}
	}
	return result
}

세 가지 결정이 이 스무 줄에 들어 있습니다.

첫째, 에러를 반환하지 않습니다. 조직이 Slack과 Discord를 함께 등록해뒀는데 Discord 웹훅이 만료됐다고 Slack 알림까지 못 받으면 실패한 채널 하나가 알림 기능 전체를 죽입니다. 그래서 루프는 계속 돕니다.

다만 에러를 반환하지 않는 코드는 실패를 삼키는 코드가 되기 쉽습니다. 그래서 두 가지를 함께 둡니다. 실패한 채널 타입과 조직 ID를 로그에 남기고, 성공과 실패 건수를 SendResult로 반환합니다. 호출자는 이 값을 지표로 올립니다. 흐름은 이어가되 건수는 남깁니다. “계속 진행한다”와 “없던 일로 한다”는 다릅니다.

이 반환값은 처음부터 있던 게 아닙니다. 그전에는 라우터가 아무것도 반환하지 않았고 호출자는 “보냈다”를 성공으로 세고 있었습니다. 실패해도 성공 카운터가 올라가는 지표였습니다. 관측되지 않는 실패보다 나쁜 것은 성공으로 관측되는 실패입니다.

둘째, 알 수 없는 채널 타입은 로그를 남기고 건너뜁니다. continue이므로 실패로도 세지 않습니다. 설정에 오타가 있거나 아직 지원하지 않는 타입이 들어와도 그 조직의 다른 알림은 정상 동작합니다. 실패 카운터를 올리지 않는 것은 “보내려다 실패한 것”과 “보낼 수 없는 것”을 구분하기 위해서입니다.

셋째, strings.ToUpper로 타입을 정규화합니다. 이것도 사고 뒤에 들어왔습니다. 백엔드 Prisma enum은 SLACK, DISCORD, PUSH인데 중간 경로 어딘가에서 소문자가 섞이면 전부 default 가지로 떨어져 조용히 스킵됩니다. 알림이 안 오는데 에러도 안 나는 상태입니다. 지금은 테스트가 이 경우를 고정하고 있습니다.

Go
channels := []configloader.AlertChannelConfig{
	{Type: "slack", WebhookURL: "https://example.test/hook-a"},  // lowercase
	{Type: "SLACK", WebhookURL: "https://example.test/hook-b"},  // uppercase
	{Type: "Slack", WebhookURL: "https://example.test/hook-c"},  // mixed
}
result := r.SendAll(context.Background(), 1, channels, msg, []string{"kw"})
// 3 sends, result.Success == 3

세 채널, 세 가지 몸통

발송기 세 개는 형태가 전부 다릅니다.

Slack은 Block Kit입니다. 섹션 블록 하나에 머리글과 본문을 개행으로 붙입니다.

Go
func buildPayload(msg model.ChatMessage, keywords []string) slackPayload {
	url := channelURL(msg.Platform, msg.ChannelID)
	label := fmt.Sprintf("[%s/%s]", escapeMrkdwn(msg.Platform), escapeMrkdwn(msg.StreamerName))
	keywordLabels := formatKeywordLabels(keywords)
	var headerText string
	if url != "" {
		headerText = fmt.Sprintf("*<%s|%s>*  `%s`", url, label, strings.Join(keywordLabels, ", "))
	} else {
		headerText = fmt.Sprintf("*%s*  `%s`", label, strings.Join(keywordLabels, ", "))
	}
	if msg.IsMelomingChannel {
		headerText = ":meloming_new: " + headerText
	}

	author := msg.Nickname
	if author == "" {
		author = "?"
	}
	bodyText := fmt.Sprintf("%s : %s", escapeMrkdwn(author), escapeMrkdwn(msg.Message))

	return slackPayload{
		Blocks: []slackBlock{
			{
				Type: "section",
				Text: &slackText{Type: "mrkdwn", Text: headerText + "\n" + bodyText},
			},
		},
	}
}

머리글의 [플랫폼/스트리머명]은 방송 페이지로 가는 링크입니다. 링크 주소는 플랫폼별로 규칙이 다릅니다.

Go
func channelURL(platform, channelID string) string {
	switch platform {
	case "chzzk":
		return "https://chzzk.naver.com/live/" + channelID
	case "soop":
		return "https://play.sooplive.com/" + channelID
	case "cime":
		return "https://ci.me/@" + channelID + "/live"
	default:
		return ""
	}
}

알림을 받는 사람이 하는 일은 “이거 무슨 맥락이지”를 확인하는 것이므로 알림에서 방송으로 가는 거리가 클릭 한 번이어야 합니다. 모르는 플랫폼이면 빈 문자열을 돌려주고 링크 없는 머리글로 대체합니다.

여기서 보안 처리가 하나 붙습니다.

Go
// escapeMrkdwn escapes Slack mrkdwn control characters in user-supplied strings.
// Prevents injection of <!channel>, <!here>, and crafted links.
func escapeMrkdwn(s string) string {
	s = strings.ReplaceAll(s, "&", "&amp;")
	s = strings.ReplaceAll(s, "<", "&lt;")
	s = strings.ReplaceAll(s, ">", "&gt;")
	return s
}

알림 본문에는 채팅 원문이 그대로 들어갑니다. 즉 임의의 시청자가 우리 Slack 워크스페이스에 문자열을 주입할 수 있는 경로입니다. <!channel>을 채팅에 치면 전 채널 멘션이 울릴 수 있고, <url|텍스트> 문법으로 가짜 링크를 만들 수도 있습니다. 외부 입력이 사내 채널까지 도달하는 경로에서는 escape가 선택이 아닙니다.

IsMelomingChannel 접두 이모지는 뒤늦게 얹은 기능입니다. 알림이 걸린 방송이 우리 서비스에 이미 연결된 채널인지 표시합니다. 판정은 백엔드 공개 조회 API 호출입니다.

Go
endpoint := fmt.Sprintf(
	"%s/v1/channel/by-platform/%s/%s",
	l.baseURL,
	url.PathEscape(platform),
	url.PathEscape(platformChannelID),
)

응답 200이면 등록, 404면 미등록, 그 외 상태 코드는 에러입니다. 알림 한 건마다 HTTP 호출을 하면 발송이 느려지므로 캐시를 답니다.

Go
const (
	registeredChannelCacheTTL   = 5 * time.Minute
	unregisteredChannelCacheTTL = time.Minute
)

긍정 결과는 5분, 부정 결과는 1분입니다. 비대칭인 이유는 상태 전이 방향이 비대칭이기 때문입니다. 미등록 채널이 등록되는 일은 자주 있고 그때 빨리 반영되면 좋지만 등록된 채널이 갑자기 사라지는 일은 드뭅니다. HTTP 타임아웃은 2초입니다.

조회에 실패했을 때의 동작이 중요합니다.

Go
registered, err := n.channelLookup.IsRegistered(context.Background(), msg.Platform, msg.ChannelID)
if err != nil {
	log.Printf("[slack] meloming channel lookup failed (platform=%s channel_id=%s): %v", ...)
} else {
	msg.IsMelomingChannel = registered
}

로그만 남기고 알림은 그대로 보냅니다. 이모지 하나 때문에 알림 자체를 막을 이유는 없습니다. 부가 정보의 실패가 본체의 실패로 번지지 않도록 경계를 그었습니다. 테스트도 이 동작을 고정합니다(TestSlackNotifier_lookupFailureStillSendsWithoutEmoji).

Discord는 embed입니다.

Go
type discordEmbed struct {
	Title       string `json:"title"`
	Description string `json:"description"`
	Color       int    `json:"color"`
	Timestamp   string `json:"timestamp"`
}

제목은 “키워드 알림” 뒤에 플랫폼과 스트리머명을 붙인 한 줄이고 본문은 매칭 키워드 목록과 채팅 원문입니다.

Go
keywordLabels := formatKeywordLabels(keywords)
description := fmt.Sprintf("매칭 키워드: %s\n%s : %s",
	strings.Join(keywordLabels, ", "),
	nickname,
	msg.Message,
)

// Use the message timestamp if valid RFC3339, otherwise use current time.
ts := msg.Timestamp
if _, err := time.Parse(time.RFC3339, ts); err != nil {
	ts = time.Now().UTC().Format(time.RFC3339)
}

Discord는 embed에 잘못된 타임스탬프가 들어가면 요청 자체를 거절합니다. 스트림에서 온 값이 항상 유효한 RFC3339라는 보장이 없으므로 파싱해보고 실패하면 현재 시각으로 대체합니다. 외부 시스템에 넘기기 직전에 형식을 확인하는 자리입니다.

성공 판정도 Slack과 다릅니다. Slack은 StatusCode == 200만 성공이고 Discord는 204 No Content를 돌려주므로 2xx 범위 전체를 성공으로 봅니다.

Go
// Discord returns 204 No Content on success.
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
	return nil
}

재시도는 양쪽 다 3회, 지수 백오프 1초, 2초, 4초입니다. 초기 백오프는 필드로 빼서 테스트에서 1밀리초로 낮춥니다.

푸시는 HTTP가 아니라 gRPC입니다. FCM이나 APNs를 직접 호출하지 않습니다.

Go
// GrpcPushClient sends push notifications via gRPC to the backend's
// KeywordAlertCommandService.SendPushNotification.
func (c *GrpcPushClient) SendPush(ctx context.Context, orgID int32, title, body, url string) error {
	conn, err := grpc.NewClient(c.addr, grpc.WithTransportCredentials(insecure.NewCredentials()))
	if err != nil {
		return fmt.Errorf("grpc dial: %w", err)
	}
	defer conn.Close()

	callCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
	defer cancel()

	client := businessv1.NewKeywordAlertCommandServiceClient(conn)
	_, err = client.SendPushNotification(callCtx, &businessv1.SendPushNotificationRequest{
		OrganizationId: orgID,
		Title:          title,
		Body:           body,
		Url:            url,
	})
	...
}

이 서비스가 아는 것은 조직 ID뿐입니다. 그 조직에 누가 속해 있고 어떤 기기 토큰을 갖고 있는지는 전혀 모릅니다. 수신자 해석은 백엔드가 합니다.

TypeScript
@GrpcMethod('KeywordAlertCommandService', 'SendPushNotification')
async sendPushNotification(data: {...}) {
  const members = await this.prisma.organizationMember.findMany({
    where: { organizationId: data.organization_id },
    select: { userId: true },
  });

  for (const member of members) {
    await this.notifications.enqueueSend(member.userId, {
      type: 'KEYWORD_ALERT',
      title: data.title,
      body: data.body,
      url: data.url,
    });
  }

  return {};
}

조직 멤버 전원에게 KEYWORD_ALERT 타입 알림을 큐에 넣습니다. 푸시 토큰, 기기 관리, 재시도, 구독 만료 처리는 이미 백엔드 알림 모듈이 하고 있는 일이므로 그대로 씁니다. 웹훅은 URL 하나만 알면 되지만 푸시는 사람과 기기를 알아야 합니다. 그 지식을 이 서비스가 가질 이유는 없습니다.

공통 표기 규칙이 하나 더 있습니다. 정규식 키워드를 알림에 그대로 찍으면 본문이 지저분해집니다. 그래서 정규식 지원에 이어 라벨 축약이 들어왔습니다.

Go
func formatKeywordLabel(keyword string) string {
	trimmed := strings.TrimSpace(keyword)
	if !isRegexKeyword(trimmed) {
		return keyword
	}
	return fmt.Sprintf("regex(%d)", keywordTraceID(trimmed))
}

func keywordTraceID(keyword string) uint32 {
	h := fnv.New32a()
	_, _ = h.Write([]byte(keyword))
	return h.Sum32()
}

일반 키워드는 그대로 보여주고, 정규식은 짧은 식별자로 바꿉니다. /hello\s+meloming/iregex(1538449266)이 됩니다. FNV-1a 해시이므로 같은 패턴은 항상 같은 숫자가 되고, 설정 화면과 대조해서 어떤 패턴이 걸렸는지 역추적할 수 있습니다. 테스트가 이 안정성을 검사합니다(TestFormatKeywordLabel_isStable).

축약이 표시 계층에서만 일어난다는 점이 중요합니다. ClickHouse에는 원본 패턴 문자열이 그대로 들어갑니다. 사람이 읽는 자리에서는 줄이고, 기계가 검색하는 자리에서는 원본을 남깁니다.


9. ACK를 가장 먼저 등록합니다

processMessage의 첫 부분입니다.

Go
func (c *Consumer) processMessage(ctx context.Context, msg redis.XMessage) {
	c.metrics.MessagesConsumed.Add(1)
	c.consumedCount.Add(1)

	// ACK with an independent context so SIGTERM cancellation does not skip the ACK.
	// A 5s deadline prevents a hung Redis connection from blocking shutdown indefinitely.
	// Registered first (before any early return) to avoid PEL accumulation on excluded senders.
	ackCtx, ackCancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer ackCancel()
	defer func() {
		if err := c.client.XAck(ackCtx, streamKey, c.group, msg.ID).Err(); err != nil {
			log.Printf("consumer: XACK failed for %s: %v", msg.ID, err)
		}
	}()

이 다섯 줄에 세 가지 결정이 들어 있습니다.

독립 컨텍스트입니다. 부모 ctx를 쓰면 SIGTERM으로 취소된 순간 ACK도 취소됩니다. 이미 처리한 메시지가 PEL(Pending Entries List)에 남아 다음 기동 때 다시 배달됩니다. 종료 신호는 “새 일을 그만 받아라”이지 “하던 일을 지워라”가 아닙니다.

5초 데드라인입니다. 독립 컨텍스트는 취소되지 않으므로 Redis가 멈춰 있으면 종료가 영원히 안 끝날 수 있습니다. 취소로부터 자유롭게 만든 컨텍스트에는 반드시 자체 시한을 답니다.

가장 먼저 등록합니다. 이 defer 아래에는 이른 반환이 네 개 있습니다.

Go
	chatMsg := parseMessage(msg.Values)

	if chatMsg.Type == "system" {
		return
	}
	if strings.EqualFold(chatMsg.Platform, "chzzk") && chatMsg.UserID == "SYSTEM_MESSAGE" {
		return
	}
	if chatMsg.Message == "" {
		return
	}
	if c.isExcludedSender(chatMsg.Platform, chatMsg.Nickname) {
		return
	}

시스템 공지, 플랫폼 시스템 계정, 빈 메시지, 제외 발신자입니다. 이 필터가 defer보다 위에 있으면 걸러진 메시지가 ACK 없이 빠져나가 PEL에 쌓입니다. 알림도 안 오고 스트림 상태만 나빠지는 조용한 누수입니다. 그래서 ACK 등록을 어떤 이른 반환보다도 앞에 뒀습니다.

제외 발신자는 채팅봇 대응입니다. 방송에 붙은 봇이 브랜드명을 반복 출력하면 알림이 봇 소리로 가득 찹니다. 형식은 platform/nickname 쌍의 쉼표 구분 목록이고 비교는 양쪽 다 소문자로 바꾼 뒤에 합니다.

Go
func (c *Consumer) isExcludedSender(platform, nickname string) bool {
	platformLower := strings.ToLower(platform)
	nicknameLower := strings.ToLower(nickname)
	for _, ex := range c.excludedSenders {
		if strings.ToLower(ex.Platform) == platformLower && strings.ToLower(ex.Nickname) == nicknameLower {
			return true
		}
	}
	return false
}

시스템 메시지 필터는 뒤에 덧붙은 방어선입니다. 플랫폼 공지와 커넥터 이벤트는 채팅 수집 서비스가 이미 type=system으로 표시해서 보내주는데도, 같은 판정을 여기서 한 번 더 합니다.

상류에서 이미 막고 있는 것을 하류에서 또 막습니다. 중복이지만 두는 이유는 배포 시점이 다르기 때문입니다. 상류의 구버전 워커가 아직 돌고 있으면 하류의 필터만이 유일한 방어선입니다. 여러 서비스가 각자의 속도로 배포되는 환경에서는 상류가 고쳐졌다는 사실과 지금 이 순간 상류가 고쳐져 있다는 사실이 다릅니다.

한편 ACK 우선 정책에는 대가도 있습니다. ClickHouse 적재와 알림 발송은 부모 ctx를 쓰므로 종료 직전 처리 중이던 메시지는 알림이나 이력 기록이 취소된 채로 ACK될 수 있습니다. 무한 재배달 루프를 피하는 대신 극히 드문 유실을 받아들인 선택입니다. 그래서 이 트레이드오프는 TestRun_acksEvenWhenNotifyFails라는 테스트로 고정해뒀습니다. 의도한 동작이라는 사실이 실행 가능한 형태로 남아 있지 않으면 나중에 누군가 이 동작을 버그로 보고 “고치게” 됩니다.


10. ClickHouse: 무엇을, 몇 행으로 남기는가

이력 적재는 배치 INSERT입니다.

Go
// MatchRecord represents a single keyword match to be written to ClickHouse.
type MatchRecord struct {
	OrganizationID int32
	Keyword        string
	StreamID       string // Redis Stream message ID
	Msg            model.ChatMessage
	MatchedAt      time.Time
}

func (w *ClickHouseWriter) Write(ctx context.Context, records []MatchRecord) error {
	if len(records) == 0 {
		return nil
	}

	batch, err := w.conn.PrepareBatch(ctx, `INSERT INTO keyword_alert_matches
		(organization_id, keyword, stream_id, platform, channel_id, streamer_name,
		 user_id, nickname, message, matched_at)`)
	...
}

테이블 정의는 다음과 같습니다.

SQL
-- ClickHouse DDL
keyword_alert_matches (
    organization_id  UInt32,
    keyword          String,
    stream_id        String,                  -- Redis Stream message ID
    platform         LowCardinality(String),  -- "chzzk", "soop", "cime"
    channel_id       String,
    streamer_name    String,
    user_id          String,
    nickname         String,
    message          String,
    matched_at       DateTime('Asia/Seoul'),
    inserted_at      DateTime DEFAULT now()
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(matched_at)
ORDER BY (organization_id, matched_at, stream_id)
TTL matched_at + INTERVAL 90 DAY;

행 수 계산이 중요합니다.

Go
for orgID, res := range orgKeywordResults {
	for _, kw := range res.labels {
		records = append(records, store.MatchRecord{
			OrganizationID: orgID,
			Keyword:        kw,
			StreamID:       msg.ID,
			Msg:            chatMsg,
			MatchedAt:      now,
		})
	}
}

채팅 한 줄이 두 조직에 걸리고 각 조직에서 키워드 두 개씩 걸리면 4행입니다. 메시지 본문이 4번 중복 저장됩니다. 정규화된 관계형 사고로는 낭비지만 컬럼 지향 저장소에서 같은 문자열이 반복되는 비용은 압축이 상당 부분 흡수합니다. 대신 조회가 단순해집니다. 조회 축이 (조직, 기간, 키워드)이므로 정렬 키가 그대로 그 축입니다.

stream_id를 남기는 이유는 재처리 대비입니다. Redis Stream 메시지 ID가 그대로 들어가므로 같은 메시지를 두 번 처리했는지 사후에 판별할 수 있습니다. ACK 우선 정책 덕분에 실제 중복은 거의 없지만 판별 수단 자체를 없애지는 않았습니다.

세 가지를 더 짚습니다.

적재 대상은 키워드 매칭뿐입니다. 사용자 지목 매칭은 알림만 가고 이력에는 남지 않습니다. keyword 컬럼이 NOT NULL인 스키마에 그룹 이름을 넣기 시작하면 “이 값이 키워드인가 그룹 이름인가”를 조회하는 쪽이 매번 판별해야 합니다. 스키마를 건드리는 대신 적재 범위를 좁게 유지했습니다.

저장하는 키워드는 축약 전 원본입니다. 8장에서 본 regex(...) 라벨은 알림 본문에만 적용됩니다.

적재 실패는 알림을 막지 않습니다.

Go
if err := c.store.Write(ctx, records); err != nil {
	log.Printf("consumer: ClickHouse write error: %v", err)
}

로그만 남기고 계속 진행합니다. 이력은 사후 분석용이고 알림은 실시간 목적이므로 둘 중 하나가 죽어야 한다면 이력 쪽입니다. 다만 6장에서 본 것처럼 기동 시점의 연결 실패는 fail-closed입니다. “쓸 수 있는데 한 건 실패”와 “처음부터 못 쓴다”를 다르게 다룹니다.

읽기 쪽은 백엔드가 담당합니다.

TypeScript
const conditions = ['organization_id = {orgId:UInt32}'];
if (query.keyword)  conditions.push('keyword = {keyword:String}');
if (query.platform) conditions.push('platform = {platform:String}');
if (query.from)     conditions.push("matched_at >= {from:DateTime64(3, 'UTC')}");
if (query.to)       conditions.push("matched_at <= {to:DateTime64(3, 'UTC')}");

조건이 전부 정렬 키와 파티션 키 위에 있습니다. 쓰기 쪽 스키마를 읽기 쪽 질의 모양에 맞춰 정한 결과입니다.


11. 지표 이름에 남은 단일 테넌트 시절의 흔적

지표는 네 개입니다.

Go
MessagesConsumed: prometheus.NewCounter(prometheus.CounterOpts{
	Name: "keyword_alert_messages_consumed_total",
	Help: "Total chat messages consumed from chat:firehose",
}),
Matches: prometheus.NewCounter(prometheus.CounterOpts{
	Name: "keyword_alert_matches_total",
	Help: "Total messages matching at least one keyword",
}),
SlackSent: prometheus.NewCounterVec(prometheus.CounterOpts{
	Name: "keyword_alert_slack_sent_total",
	Help: "Total Slack notifications attempted",
}, []string{"result"}),
SlackDuration: prometheus.NewHistogram(prometheus.HistogramOpts{
	Name:    "keyword_alert_slack_send_duration_seconds",
	Help:    "Slack webhook send duration in seconds",
	Buckets: prometheus.DefBuckets,
}),

레이블은 SlackSentresult 하나이고 값은 successerror입니다.

그런데 실제 사용처를 보면 이름과 내용이 어긋나 있습니다.

Go
start := time.Now()
totalSuccess, totalFailure := 0, 0
for orgID, res := range orgKeywordResults {
	result := c.router.SendAll(ctx, orgID, res.channels, chatMsg, res.labels)
	totalSuccess += result.Success
	totalFailure += result.Failure
}
for orgID, res := range orgUserResults {
	result := c.router.SendAll(ctx, orgID, res.channels, chatMsg, res.labels)
	...
}
c.metrics.SlackDuration.Observe(time.Since(start).Seconds())

if totalSuccess > 0 {
	c.metrics.SlackSent.WithLabelValues("success").Add(float64(totalSuccess))
}
if totalFailure > 0 {
	c.metrics.SlackSent.WithLabelValues("error").Add(float64(totalFailure))
}

keyword_alert_slack_sent_total은 Slack만이 아니라 Discord와 푸시까지 포함한 전 채널 발송 건수입니다. keyword_alert_slack_send_duration_seconds는 한 메시지에 대한 모든 조직, 모든 채널 발송의 합계 시간입니다. Slack밖에 없던 단일 테넌트 시절의 이름이 그대로 남았습니다.

이름을 바꾸지 않은 이유는 대시보드와 알럿 룰이 이미 이 이름을 참조하고 있었기 때문입니다. 지표 이름 변경은 코드 한 줄이지만 관측 스택 전체의 변경입니다. 다만 대가가 있습니다. 이 지표만 보고는 어느 채널이 실패했는지 알 수 없습니다. 채널별 분해가 필요하면 result 옆에 channel_type 레이블을 추가하는 것이 정공법이고 그때가 이름을 함께 정리할 시점입니다.

이름과 내용이 어긋난 지표는 다음 사람에게 버그로 보입니다. 의도적으로 남긴 부채는 의도적으로 표시해 두어야 부채로 남습니다.

지표 서버와 헬스 서버는 포트를 나눠 띄웁니다.

Go
healthSrv := &http.Server{Addr: ":8080", Handler: healthMux}   // /healthz
metricsSrv := &http.Server{Addr: ":9090", Handler: metricsMux} // /metrics

1분마다 처리량도 로그로 남깁니다.

Go
func (c *Consumer) logStats(ctx context.Context) {
	ticker := time.NewTicker(1 * time.Minute)
	...
	current := c.consumedCount.Load()
	log.Printf("stats: consumed %d messages in the last minute (total: %d)", current-last, current)
}

Prometheus가 있는데도 로그를 남기는 것은 파드 로그만 보고도 “지금 흐르고 있는가”를 즉시 확인하기 위해서입니다. 카운터는 atomic.Uint64입니다.


12. 테스트가 지키는 경계

테스트는 13개 파일에 75개 함수입니다. 무엇을 고정해뒀는지가 곧 무엇을 깨뜨리기 쉬운지의 목록입니다.

매처(13개)

  • TestMatch_caseInsensitive: meloming 키워드가 Hello MELOMING world에 걸리고, 결과 라벨은 원본 표기 그대로입니다.
  • TestMatch_regexPrefix: regex:신청(곡|합니다)노래 신청합니다에 걸립니다.
  • TestMatch_slashRegexWithFlags: /hello\s+meloming/iHELLO MELOMING에 걸립니다. 슬래시 파싱, 플래그 변환, RE2 인라인 문법 조립이 한 번에 검증됩니다.
  • TestMultiTenantMatch: 조직 3개에 서로 다른 키워드를 주고 한 메시지를 넣어, 걸린 조직 2개만 결과에 있고 나머지 하나는 키 자체가 없음을 확인합니다.
  • TestMatchGroups_returnsOnlyMatchingGroups: 두 그룹 중 걸린 그룹만 반환되고, 그 그룹의 매칭 키워드가 2개임을 확인합니다.
  • TestMatchGroups_emptyMessage_noMatch: 빈 메시지가 어떤 키워드에도 걸리지 않아야 합니다. strings.Contains(x, "")가 항상 참이라는 Go의 성질 때문에 빈 키워드 필터링과 함께 반드시 고정해야 하는 경계입니다.

중복 제거(2개)

  • TestDedupeGroupMatches_sameChannelInTwoGroups: 그룹 A가 채널 1, 2를, 그룹 B가 채널 1, 3을 가리킬 때 결과가 채널 3개와 키워드 2개인지 확인합니다. 7장의 정책이 그대로 테스트 이름이 됐습니다.

라우터(6개)

  • TestRouter_SendAll_errorsDoNotBlockOthers: Slack이 에러를 반환해도 Discord와 푸시가 호출됐는지 확인합니다.
  • TestRouter_SendAll_unknownTypeSkipped: email 타입이 섞여 있어도 뒤의 Slack 채널이 발송됩니다.
  • TestRouter_SendAll_CaseInsensitive: slack, SLACK, Slack 세 표기가 모두 Slack으로 라우팅되고 성공 카운트가 3입니다.

컨슈머(9개)

miniredis로 실제 Redis Streams 동작을 흉내 냅니다.

  • TestRun_acksEvenWhenNotifyFails: 처리 후 XPending 카운트가 0인지 확인합니다. PEL이 비어 있다는 것이 ACK 되었다는 증거입니다.
  • TestRun_dropsSystemAndEmptyMessages: type=system, userId=SYSTEM_MESSAGE, 빈 메시지 세 건을 넣고 정상 메시지 한 건을 더 넣어, 발송이 정확히 한 번인지 확인합니다.
  • TestRun_excludesSenders: 제외 목록에 있는 닉네임의 메시지는 걸러지고 다른 시청자의 같은 문장은 통과합니다.
  • TestRun_gracefulShutdown: 컨텍스트 취소 후 2초 안에 Run()이 nil을 반환하는지 확인합니다.
  • TestRun_noConfigsSkipsNotification: 설정이 비어 있으면 소비는 하되 발송은 하지 않습니다.

설정 로더(4개)와 노티파이어(24개)

  • TestGetConfigs_deepCopy: 반환된 스냅샷을 수정해도 로더 내부가 변하지 않는지 확인합니다. 6장의 동시성 전제가 이 테스트 하나에 걸려 있습니다.
  • TestSlackNotifier_retriesOnServerError: 500을 두 번 돌려주다 세 번째에 200을 주면 성공이고 시도 횟수가 정확히 3입니다.
  • TestHTTPMelomingChannelLookup_registeredAndCached: 같은 채널을 두 번 조회했을 때 실제 HTTP 요청이 1회인지 확인합니다.
  • TestBuildDiscordPayload_invalidTimestamp: 잘못된 타임스탬프가 현재 시각으로 대체되는지 확인합니다.

환경변수 파싱(15개)과 ClickHouse 레코드(2개)

나머지 17개는 환경변수 로딩과 적재 레코드 구성입니다. 특히 TestLoad_emptyKeywordsAreFiltered는 빈 키워드가 걸러지는지 확인합니다. strings.Contains(x, "")가 항상 참이므로 빈 문자열 키워드 하나가 남으면 그 조직은 모든 채팅에 알림을 받게 됩니다. TestLoad_excludedSendersIgnoresInvalidFormatplatform/nickname 형식이 깨진 항목을 건너뛰는지 확인합니다.

CI는 go test -v -race ./...입니다. 레이스 디텍터를 켜기 때문에 CGO가 필요하고 러너에 C 툴체인을 설치하는 단계가 워크플로에 들어 있습니다. 인메모리 스냅샷을 여러 고루틴이 읽는 구조라면 레이스 검출을 빼고 갈 수 없습니다.


13. 배포 구성

빌드는 2단계 Dockerfile입니다.

Dockerfile
FROM golang:1.26-alpine AS builder
ENV GOPRIVATE=github.com/<org>/*
...
RUN --mount=type=secret,id=github_token \
    if [ -f /run/secrets/github_token ]; then \
      git config --global url."https://x-access-token:$(cat /run/secrets/github_token)@github.com/".insteadOf "https://github.com/"; \
    fi && \
    go mod download && ...

RUN CGO_ENABLED=0 GOOS=linux go build -o /keyword-alert-service ./cmd/

FROM gcr.io/distroless/static-debian12:nonroot
COPY --from=builder /keyword-alert-service /keyword-alert-service
USER nonroot:nonroot
ENTRYPOINT ["/keyword-alert-service"]

공용 proto 모듈이 private 저장소라 토큰이 필요한데 --build-arg로 넘기면 이미지 레이어 히스토리에 그대로 남습니다. 그래서 BuildKit secret 마운트로 바꿨습니다.

런타임은 distroless static입니다. 셸도 패키지 매니저도 없고 non-root로 뜹니다. 정적 링크 Go 바이너리 하나만 있으면 되는 워크로드에 리눅스 배포판 전체를 얹을 이유가 없습니다.

배포는 GitHub Actions에서 이미지를 만들어 사내 컨테이너 레지스트리에 올린 뒤 배포 설정 저장소의 values 파일에서 태그를 바꿔 커밋하는 방식입니다. qa 브랜치는 QA 환경으로, main 브랜치는 프로덕션으로 갑니다. 클러스터에 직접 kubectl apply 하는 단계는 없습니다.

YAML
# .github/workflows/deploy-prod.yml (요약)
env:
  IMAGE_REPOSITORY: <사내 컨테이너 레지스트리>/keyword-alert-service

jobs:
  test:            # go test -race + build check
  build-and-push:  # runs-on: ubuntu-24.04-arm, --platform linux/arm64
  deploy:          # 배포 설정 저장소의 프로덕션 values 파일 태그 갱신 후 push

프로덕션 이미지는 arm64입니다. 노드 셀렉터가 그렇게 요구합니다.

YAML
replicaCount: 1
revisionHistoryLimit: 1

strategy:
  type: Recreate

env:
  - name: CONSUMER_NAME
    valueFrom:
      fieldRef:
        fieldPath: metadata.name
  - name: EXCLUDED_SENDERS
    value: "chzzk/example-bot,soop/example-bot"
  - name: GRPC_ADDR
    value: "<백엔드 gRPC 주소>"
  - name: REDIS_ADDR
    value: "<채팅 버퍼 Redis 주소>:6379"
  - name: CLICKHOUSE_ADDR
    value: "<채팅 버퍼 ClickHouse 주소>:9000"

resources:
  # Right-sized from 30d peak: mem 16Mi, cpu small.
  requests: { cpu: 50m, memory: 64Mi }
  limits:   { cpu: 300m, memory: 192Mi }

service:     { enabled: false }
ingress:     { enabled: false }
autoscaling: { enabled: false }

nodeSelector:
  kubernetes.io/arch: arm64

tolerations:
  - key: node.kubernetes.io/lifecycle
    operator: Equal
    value: spot
    effect: NoSchedule

몇 가지가 서로 맞물려 있습니다.

replicaCount: 1strategy: Recreate. 컨슈머 그룹을 쓰므로 원리상 여러 파드가 나눠 읽을 수 있습니다. 그럼에도 하나로 두고 롤링이 아닌 재생성 방식으로 배포합니다. 30일 피크 기준 메모리 16Mi인 워크로드에 분산 처리가 필요 없고 배포 중 두 파드가 겹쳐 뜨는 순간을 아예 없애는 편이 단순하기 때문입니다.

CONSUMER_NAME은 파드 이름입니다. fieldRefmetadata.name을 넣습니다. 파드가 교체될 때마다 컨슈머 이름이 바뀝니다. ACK를 항상 하기 때문에 옛 이름의 PEL이 비어 있습니다. 그래서 이름이 바뀌어도 고아 항목이 남지 않습니다. 9장의 ACK 정책과 여기가 연결됩니다.

service: false, ingress: false. 이 서비스는 아무 요청도 받지 않습니다. 들어오는 것은 스트림에서 당겨온 메시지뿐이고, 나가는 것은 웹훅과 gRPC입니다. 8080과 9090 포트는 클러스터 내부의 프로브와 스크레이프 용도입니다.

비밀값은 시크릿 오퍼레이터가 외부 비밀 저장소에서 읽어 파드 환경변수로 주입합니다.

YAML
secrets:
  enabled: true
  refreshInterval: 1m
  keys:
    - SLACK_WEBHOOK_URL
    - CONSUMER_GROUP
    - CLICKHOUSE_PASSWORD

CONSUMER_GROUP이 비밀값 목록에 있는 것이 눈에 띕니다. 컨슈머 그룹 이름은 비밀이 아니지만 같은 비밀 저장소에서 함께 관리하면 값 하나를 바꿀 때 참조 위치가 하나로 유지됩니다. 그룹 이름을 바꾸는 것은 “이 스트림을 처음부터 다시 읽겠다”에 가까운 조작이라 아무나 values 파일에서 고치지 않게 하는 효과도 있습니다.

관측은 프로메테우스 스크레이프 대상으로 등록해서 붙입니다.

YAML
metrics:
  enabled: true
  port: 9090
  path: /metrics
  scrape:
    enabled: true
    interval: 30s

같은 서비스가 온프렘 쿠버네티스 클러스터에도 배포됩니다. 값 파일이 클라우드용과 온프렘용으로 나뉘고 차이는 노드 셀렉터와 의존 서비스 주소뿐입니다. 온프렘 쪽은 온프렘에 있는 채팅 버퍼 Redis와 ClickHouse를 보고 gRPC는 사설망 주소로 넘어갑니다. 애플리케이션 코드는 한 줄도 다르지 않습니다. 의존 대상 주소를 전부 환경변수로 뺀 덕분에 인프라 이전이 값 파일 문제로 남았습니다.


14. 마무리

이 글은 키워드 알림 서비스가 조직, 그룹, 채널 3층 구조로 자라는 동안 어디를 고치고 어디를 그대로 뒀는지 따라간 기록입니다. 스트림 읽기는 하나로 두고 매처만 멀티테넌트로 만들었고, 메시지마다 설정 스냅샷을 새로 받는 구조가 설정 갱신 경합과 정규식 재컴파일 비용을 함께 정리했습니다.

판정과 전달의 경계도 그 과정에서 정해졌습니다. 매처는 무엇이 걸렸는지만 답하고 그룹 사이 중복 제거와 채널별 발송은 호출자와 라우터가 맡습니다. 한 채널이 실패해도 나머지는 나가되 실패 건수는 지표에 남고, 설정을 못 받거나 ClickHouse에 붙지 못하면 아예 뜨지 않습니다. ACK를 어떤 이른 반환보다 먼저 등록한 것도 같은 판단입니다.

입력이 되는 채팅 스트림 자체를 끊기지 않게 유지하는 방법은 SOOP, 치지직, 씨미 채팅 수집 커넥터로 12,000개 방송 채팅 수집하기에 정리했습니다.

시리즈 · 채팅 수집과 처리

4 / 5

SOOP, 치지직, 씨미의 라이브 채팅을 모아 랭킹과 알림으로 바꾸기까지의 기록.

  1. 1.멜로밍 랭킹 - SOOP, 치지직, 씨미 채팅 수집 및 처리 아키텍처
  2. 2.SOOP, 치지직, 씨미 채팅 수집 커넥터로 12,000개 방송 채팅 수집하기
  3. 3.라이브 스트리밍 채팅 수집 시스템을 KEDA와 Karpenter로 확장한 방법
  4. 4.채팅 키워드알림 서비스 - 멀티테넌트 키워드 알림
  5. 5.방송에서 가장 많이 나온 단어는? : NLP를 활용한 채팅 워드클라우드