🚀 Введение

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

Поток данных:

  1. go-ingestion забирает контент из источников (RSS/API/SSE/t.me) → публикует в JetStream raw.ingest.<type>.<id>
  2. go-nlp-processor потребляет сообщения из очереди, очищает текст, кластеризует фразы → шлёт тренды по gRPC в state-aggregator
  3. go-state-aggregator хранит тренды в памяти, применяет физику, публикует кадры в frames.live (~каждые 500 мс) и пишет бизнес-данные (пользователи, источники) в SQLite
  4. go-websocket-edge — stateless-шлюз: подписан на frames.live, раздаёт кадры браузерам
  5. 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 до продакшена.

Буду рад вопросам и идеям!