🚀 Введение
VibeRadar — это живой радар трендов. Он подключается к десяткам потоков текста (Wikipedia EventStreams, RSS-ленты, Telegram-каналы), превращает их в набор ключевых фраз и рисует их на карте как пузыри: чем громче тема — тем больше пузырь. И всё это в реальном времени.
Проект родился из простого вопроса: а что будет, если взять шум интернета и попытаться увидеть в нём зарождающиеся тренды за секунды, а не за часы? Ответ — распределённый бэкенд на Go из пяти микросервисов, очередь на NATS JetStream, NLP-кластеризация без внешних LLM и «физика» для анимированных пузырей.
В статье разберу внутренности: как очищается текст, как фразы кластеризуются в тренды, как работает семантическое слияние парафраз и почему я гонялся за каждой аллокацией в горячих путях.
🏗️ Архитектура
Пять Go-сервисов в одном модуле, protobuf/gRPC между ними, NATS JetStream как транспортная шина:
flowchart LR
subgraph Sources["Источники"]
RSS["RSS / API"]
SSE["SSE (Wikipedia)"]
TME["Telegram (t.me)"]
end
subgraph Dock["Docker Network"]
ING["go-ingestion\n+ go-discovery-worker\n(аналитик)"]
JS[("NATS JetStream\nraw.ingest.*")]
NLP["go-nlp-processor"]
EMB["Embedder (Python)"]
ST["go-state-aggregator"]
EDGE["go-websocket-edge"]
end
BR[("Браузер")]
RSS --> ING
SSE --> ING
TME --> ING
ING -->|"publish"| JS
JS -->|"consume"| NLP
NLP <-->|"gRPC"| ST
NLP -->|"эмбеддинги"| EMB
EMB -->|"векторы"| NLP
ST -->|"frames.live (500 мс)"| EDGE
EDGE <-->|"WSS"| BR
Поток данных:
- go-ingestion забирает контент из источников (RSS/API/SSE/t.me) → публикует в JetStream
raw.ingest.<type>.<id> - go-nlp-processor потребляет сообщения из очереди, очищает текст, кластеризует фразы → шлёт тренды по gRPC в state-aggregator
- go-state-aggregator хранит тренды в памяти, применяет физику, публикует кадры в
frames.live(~каждые 500 мс) и пишет бизнес-данные (пользователи, источники) в SQLite - go-websocket-edge — stateless-шлюз: подписан на
frames.live, раздаёт кадры браузерам - go-discovery-worker — аналитик источников: анализирует контент, выявляет и проверяет новые источники
📡 Источники
VibeRadar питается реальными потоками: SSE-стрим Wikipedia (ru), RSS-ленты ТАСС, РБК, РИА Новости, Интерфакса и Медузы, Habr, vc.ru, dtf.ru и Telegram-канал breakingmash. Mock-источники остались как возможность для демо и офлайн-разработки.
Источник — это конфиг, а не код: тип подключения, URL, периодичность опроса и JSON-правила фильтрации и извлечения. Пример — Wikipedia (ru), где фильтр отсекает ботов и служебные правки, а язык определяется автоматически по домену:
{
"id": "ru-wiki",
"type": "sse",
"url": "https://stream.wikimedia.org/v2/stream/recentchange",
"interval_sec": 30,
"config": {
"filter": {
"meta.domain": { "$eq": "ru.wikipedia.org" },
"bot": { "$ne": "true" },
"type": { "$in": ["edit", "new"] }
},
"extract": {
"text": ["title"],
"meta": { "lang": "auto", "url": ">title_url" }
}
}
}
Фильтры поддерживают $eq, $ne, $in, $exists, $contains и $contains_none — так, у breakingmash из постов вырезаются рекламные маркеры («#реклама», «промокод» и др.). Изменения применяются без рестарта: через REST API с JWT-авторизацией, правкой в SQLite или публикацией в NATS sources.add / sources.remove.
🕵️ Прокси и мимикрия
Часть источников — крупные СМИ и t.me — прикрыта антибот-защитой. Чтобы не светить один IP и не выглядеть ботом, у ingestion два слоя: ротационные прокси и мимикрия браузера.
Ротация прокси
Прокси задаются одной переменной PROXY_POOL — списком через запятую (http://, https://, socks5://, с логином/паролем в URL). Пустой пул — прямое подключение. Ротация — round-robin через atomic.Uint64:
func (p *ProxyPool) Transport() *http.Transport {
if len(p.urls) == 0 { return p.directT }
idx := p.index.Add(1)
return p.transports[int(idx-1)%len(p.transports)]
}
Транспорт для каждого прокси строится один раз и переиспользуется — Keep-Alive живёт между запросами, без повторных TCP+TLS рукопожатий. MaxIdleConns=1000, MaxIdleConnsPerHost=100, чтобы сотня воркеров не исчерпала сокеты (TIME_WAIT).
Мимикрия браузера
У каждого источника свой Mimicry с RNG, засеянным от ID источника. Перед запросом случайно выбирается «браузер» из пула — Chrome 131, Firefox 133, Safari 18.2, Chrome/Linux, Safari/iPhone — и под него ставятся заголовки (UA, Accept, Accept-Language, Accept-Encoding gzip, deflate, br), плюс PreDelay 50–300 мс, чтобы запросы не были роботизированными.
uTLS — согласованный TLS-отпечаток
Стандартный net/http оставляет характерный ClientHello, по которому боты палятся на уровне TLS ещё до HTTP. Рукопожатия идут через uTLS (включается UTLS_ENABLED), а отпечаток выбирается в пару к User-Agent — Chrome UA → HelloChrome_Auto, Firefox → HelloFirefox_Auto, Safari/iPhone → HelloIOS_Auto:
func (m *Mimicry) Fingerprint() utls.ClientHelloID {
ua := strings.ToLower(m.lastUA)
switch {
case strings.Contains(ua, "chrome"):
return utls.HelloChrome_Auto
case strings.Contains(ua, "firefox"):
return utls.HelloFirefox_Auto
case strings.Contains(ua, "safari"), strings.Contains(ua, "iphone"):
return utls.HelloIOS_Auto
default:
return utls.HelloRandomized
}
}
uTLS-транспорт туннелируется через прокси вручную: HTTP-прокси — CONNECT с Proxy-Authorization: Basic, SOCKS5 — отдельный dialer. Нюанс: ответ CONNECT читается байт за байтом — bufio забуферизовал бы начало TLS-стрима, и рукопожатие падало бы с tls: bad record. Поверх uTLS включается HTTP/2 (ForceAttemptHTTP2).
🛰️ Автообнаружение источников
В проде добавлять источники вручную скучно, поэтому discovery-worker работает как аналитик: он читает тот же поток raw.ingest.*, вытаскивает URL из контента, считает упоминания доменов и, когда домен набирает порог (15 упоминаний за 24 часа), проверяет его на RSS/API/SSE и предлагает как источник:
- SSRF-защита на горячем пути (проверка IP, запрет внутренних подсетей)
- blacklist доменов (6 часов) — чтобы сомнительные источники не мучили
- лимит одновременных HTTP-проверок (10), graceful shutdown
Найденный источник автоматически попадает в SQLite и через sources.add в NATS немедленно активируется в ingestion — без рестарта.
🧹 Очистка текста
Первая задача NLP-пайплайна — из сырого HTML/JSON-потока вытащить «чистый» текст. Классический подход — набор регэкспов, но это миллионы дорогих операций под нагрузкой. Вместо этого — один проход по рунам:
for _, r := range src {
if isEmoji(r) { // диапазоны, не regex
if !inSpace { out = append(out, ' '); inSpace = true }
continue
}
if unicode.IsPunct(r) || unicode.IsSymbol(r) {
if !inSpace { out = append(out, ' '); inSpace = true }
continue
}
r = unicode.ToLower(r)
if unicode.IsSpace(r) {
if !inSpace { out = append(out, ' '); inSpace = true }
continue
}
out = append(out, string(r)...)
inSpace = false
}
Плюс:
- HTML вырезается токенизатором
golang.org/x/net/html(с пропуском<script>/<style>/<svg>), а не регэкспом - URL и email вырезаются аккуратным откатом: при встрече
://код возвращается назад и удаляет схему — снова без аллокаций - Эмодзи распознаются диапазонами юникода (
0x1F000–0x1FAFFи др.), а не гигантским паттерном - Буферы из
sync.Pool, подстроки через слайсинг — ноль аллокаций в горячем пути
🧬 Кластеризация фраз
MinHash — компактный отпечаток
Из очищенного текста вытаскиваются n-граммы (2–3 слова, с русским и английским стеммингом через snowball и 64-шардовый LRU-кэш). Дальше задача — склеить дубли и вариации («ким чен ын» из двух источников) в один тренд. Полный перебор попарных сравнений слишком дорог, поэтому используется связка MinHash → LSH → точная проверка:
func (m *MinHasher) SignatureInPlace(dst []uint64, tokens []string) {
for i := range dst { dst[i] = ^uint64(0) }
for _, token := range tokens {
for i, seed := range m.seeds {
h := murmur3.SeedStringSum64(seed, token)
if h < dst[i] { dst[i] = h } // min-hash
}
}
}
LSH — кандидаты без попарного перебора
Сигнатура из 128 значений хешей сводит фразу к компактному «отпечатку». Чтобы не сравнивать каждую пару, сигнатура режется на 16 банд по 8 строк — фразы, совпавшие хотя бы в одной банде, становятся кандидатами:
func (l *LSH) BandKey(sig []uint64, b int) uint64 {
start, end := b*l.rows, b*l.rows + l.rows
h := fnv.New64a()
for i := start; i < end; i++ {
var buf [8]byte
binary.LittleEndian.PutUint64(buf[:], sig[i])
h.Write(buf[:])
}
return uint64(b)<<56 | (h.Sum64() & 0x00FFFFFFFFFFFFFF) // 0 аллокаций строк
}
Ключи банд — uint64 (индекс банды в старших 8 битах), это в ~15 раз быстрее строковых ключей. Кандидаты проверяются оценкой Жаккара по сигнатурам, и похожие фразы сливаются через Union-Find.
Вложенные фразы — там, где Jaccard не справляется
Один тонкий момент: Jaccard не видит вложенность. «Ким чен» и «ким чен ын» пересекаются на 2/3 (≈ 0.67) — ниже порога. Добивает unionContainment: если токены одной фразы — строгое подмножество другой, они склеиваются:
func isStrictSubset(a, b map[string]struct{}) bool {
if len(a) >= len(b) { return false }
for tok := range a {
if _, ok := b[tok]; !ok { return false }
}
return true
}
Каноничная фраза тренда — самая длинная, вес — сумма весов всех фраз кластера, а стабильный trend_id — 128 бит sha256(lang + phrase). Без коллизий даже на миллиардах трендов.
🫧 Физика пузырей
Самое заметное визуально — «живые» пузыри.
Затухание и демпфирование
У каждого тренда есть вес, который экспоненциально затухает: слышали про тему — вес растёт, перестали — пузырь сдувается. Затухание зависит от типа источника: у RSS/API decay = 0.99999 (медленно, тренд живёт долго), у SSE-потоков — 0.98 (новости умирают за минуты). Раз в 100 мс физический тик пересчитывает всё:
t.Weight *= math.Pow(decay, dt*10) // затухание веса
dampFactor := math.Pow(cfg.Damping, dt*10) // демпфирование скорости
t.VX *= dampFactor
t.VY *= dampFactor
t.X += t.VX // движение
t.Y += t.VY
t.VX += (0.5 - t.X) * 0.0005 // «пружина» к центру
Отталкивание и топ-300
Между пузырями работает парное отталкивание (пересеклись — разлетаются), при этом рендерится только топ-300 по весу через min-heap, чтобы O(n²) отталкивания не взорвалось на тысячах трендов. Тик сделан через time.Sleep с реальной дельтой времени, а не time.Ticker — тот под нагрузкой может «спамить» тиками, и симуляция ускоряется в разы.
🔮 Семантическое слияние
Лексическая кластеризация не умеет видеть, что «покушение на гендиректора» и «взрыв автомобиля» — это одно и то же событие, если слова разные.
Парафразы через эмбеддинги
Для этого в пайплайн добавлен Python-эмбеддер (sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2), который превращает каноничную фразу в вектор. Векторы L2-нормализованы, поэтому косинус = скалярное произведение:
func (s *inMemoryStore) semanticMatch(space map[string]*Trend, lang string, vec []float32) *Trend {
var best *Trend
var bestSim float32
for _, t := range space {
ev, ok := s.embeddings[t.TrendID]
if !ok || len(ev) != len(vec) { continue }
if sim := dot32(vec, ev); sim > bestSim {
bestSim, best = sim, t
}
}
if best != nil && float64(bestSim) >= s.cfg.SemanticMergeThreshold {
return best
}
return nil
}
Новый тренд, чьё вложение похоже на существующее выше порога (0.82 по умолчанию), не создаётся — он вливается в старый: вес суммируется, фраза сохраняется каноничная. Это чинит фрагментацию, которую не видит лексический кластеризатор.
Graceful degradation
Если эмбеддер недоступен, тренды не падают, просто семантическое слияние отключается и система деградирует до лексики. Пересечение языков (RU↔EN) выключено по умолчанию, чтобы не склеивать несвязанное.
🧱 Надёжность и производительность
Этот проект я писал в стиле «сначала продакшн-хардненинг». Из 107 коммитов добрая половина — фиксы P0 и оптимизации:
- Грациозное завершение: 20-секундный drain, workers уважают
ctx.Done(), goroutine top-level — сrecover() - Zero-alloc hot paths: publisher 4→1 аллокаций на сообщение, LSH без строковых ключей, MinHash in-place, dedup через FNV ^ seq с LRU-кэшем
- gRPC chunking: батчи по 128 трендов, чтобы не упереться в flow-control deadlock HTTP/2
- Edge на 100k соединений: stateless, JWT (HS256) проверяется локально, лимит 100 коннектов на IP, rate-limit коннектов, sharded-кэши вместо глобальных мьютексов
- OOM-защита NLP: GOMEMLIMIT, бэтчинг по 256 сообщений/256 КБ, dedup n-грамм до кластеризации
- 63 бенчмарка покрывают горячие пути: от
CleanTextдоBroadcastна 100k клиентов
Показательный кейс: SSE-подключения сначала рвались из-за короткоживущего контекста, а после фикса connectLoop — стабильно живут сутками. Именно такие вещи я считаю «настоящей» разработкой продакшена.
📊 Наблюдаемость
Пять сервисов отдают 53 метрики Prometheus (fetch rate, ошибки, глубина очередей, количество трендов, WS-соединения, семантические слияния). Их собирает vmagent и пишет в VictoriaMetrics. Из коробки — 6 дашбордов: Overview, Ingestion, NLP, State, Edge, Discovery. Например, state_semantic_merges_total и гистограмма merged_similarity помогают тюнить порог слияния по фактам, а не на глаз.
🧰 Стек
| Компонент | Роль |
|---|---|
| Go + protobuf/gRPC | пять сервисов и их общение |
| NATS JetStream | шина сообщений и очереди |
| SQLite (CGO=0) | бизнес-данные без внешней БД |
| Python + sentence-transformers | семантические эмбеддинги |
| VictoriaMetrics + Grafana | метрики и дашборды |
| Docker Compose + Traefik | упаковка и раздача |
🚀 Деплой
Всё упаковано в Docker (distroless-образы, CGO=0), Compose-стек с четырьмя режимами развёртывания:
| Режим | Когда |
|---|---|
deploy-direct | Тест/внутренняя сеть, edge на :8080 |
deploy-traefik | Встроенный Traefik + Let’s Encrypt |
deploy-traefik-ext | Уже есть внешний Traefik (мой случай) |
| Внешний nginx/caddy | Свой прокси перед edge |
Бюджет: ~3.5 ГБ RAM на стек без Traefik, 20 ГБ диска. Резервное копирование — WAL-checkpoint SQLite + снапшоты VictoriaMetrics + NATS.
💭 Заключение
VibeRadar — это эксперимент, который дорос до полноценного распределённого бэкенда: пять сервисов, очередь сообщений, NLP без облачных API и «физический» рендер данных. Самое ценное, что я вынес из проекта:
- Lexical + semantic слои дополняют друг друга: MinHash/LSH ловят 90% дублей дёшево, эмбеддинги добивают парафразы точечно
- Производительность — это про аллокации, а не про «быстрый язык»: 0-alloc пути ускоряют больше, чем смена фреймворка
- Продакшн-хардненинг — первая фича, а не последняя: каждый P0-фикс окупился десятком часов дебага в проде
Сейчас VibeRadar живёт как pet-проект с реальными источниками — от SSE Wikipedia до RSS российских СМИ и Telegram (полный список в разделе «Источники»). Mock-источники остались возможностью для демо и разработки. В планах — открыть код и довести фронтенд на Svelte 5 + PixiJS до продакшена.
Буду рад вопросам и идеям!
