Microservice: Сервис Маркетинговых Рассылок (Marketing Campaign Worker)
Спецификация изолированного Highload-воркера пакетной дистрибуции пуш-уведомлений на Go
1. Назначение и архитектурная изоляция компонента
Микросервис Marketing Campaign Worker — это специализированный изолированный воркер на языке Go, предназначенный для проведения массовых (веерных) маркетинговых рассылок на миллионы пользователей. Главная цель выделения этого компонента в отдельный сервис — защита транзакционного конвейера системы от перегрузки сети, CPU и исчерпания пулов соединений с базами данных.
Ключевые принципы интеграции и разделения ресурсов:
- Общая СУБД в режиме Read-Only: Сервис не имеет собственной базы данных. Он разделяет физическую схему таблиц справочника notification_templates и профилей пользователей с основным сервисом уведомлений. При этом маркетинговый воркер подключается строго к Read-Only реплике PostgreSQL. Запись новой информации (сегменты, пользователи, новые шаблоны) в таблицы осуществляется исключительно внешними ETL-процессами или системными миграциями.
- Изоляция очередей брокера: Маркетинговые кампании полностью выводятся из транзакционного топика. Сервис слушает выделенный топик marketing_campaigns_bulk, что исключает эффект «засорения трубы» для критически важных уведомлений пользователей.
2. Алгоритм пакетной Highload-обработки (Throttling & Multicast)
На вход в Kafka-топик сервис получает не миллион отдельных сообщений, а один облегченный пакет инициации кампании. Пакет содержит инвариантный код шаблона и вектор (массив) идентификаторов целевых пользователей.
Процесс обработки в Go-воркере организуется по следующему алгоритму:
[Kafka Message] ──► [Нарезка слайса на чанки по 500 ID]
│
┌─────────┴─────────┐
▼ ▼
[Горутина 1] [Горутина 2] ... (Worker Pool)
│ │
(Batch SELECT к Реплике БД) (Batch SELECT к Реплике БД)
│ │
[500 FCM-токенов] [500 FCM-токенов]
│ │
(messaging.SendMulticast) (messaging.SendMulticast)
│ │
└─────────┬─────────┘
▼
[Google FCM Cloud Infrastructure]
2.1. Механизм нарезки и пакетного чтения (Batching)
- Чанкование вектора: Полученный из Kafka массив user_ids (размером до 1 000 000 элементов) внутри горутины нарезается на мелкие слайсы (чанки) строго по 500 штук. Это ограничение обусловлено лимитами Google API.
- Пакетный SQL-запрос (Batch SELECT): Вместо поштучного перебора, каждая горутина-воркер выполняет к Read-Only реплике PostgreSQL один единственный запрос с оператором IN для пачки из 500 пользователей:
SELECT fcm_token FROM user_device_tokens WHERE user_id IN ($1, $2, ... $500) AND fcm_token IS NOT NULL; - Пакетная отправка (Multicast): Извлеченные 500 токенов упаковываются в структуру messaging.MulticastMessage и отправляются в Google FCM за один сетевой HTTP/2 запрос через метод messaging.SendMulticast().
2.2. Контроль пропускной способности (Rate Limiting)
Для предотвращения выжигания сетевого интернет-канала всей ноды, в рантайме Go настраивается лимитер скорости (Throttling). Воркер пул ограничивает количество одновременно отправляемых пачек в секунду (например, не более 40 пачек/сек, что эквивалентно стабильному потоку в 20 000 пушей в секунду).
3. Спецификация асинхронного контракта (Kafka Ingress)
- Целевой топик (Topic): marketing_campaigns_bulk
- Идентификатор группы (Group ID): bupar_marketing_worker_group
3.1. Структура входящего пакета (Bulk Campaign JSON)
{
"campaign_id": "camp_fall_sale_2026",
"template_code": "PROMO_DISCOUNT_30",
"x_request_id": "req-bulk-marketing-99aa",
"target_user_ids": [
"4512",
"4513",
"4514",
"7890",
"10243"
]
}4. Справочник Обсервабилити (Метрики массовых рассылок для ClickHouse)
Поскольку воркер оперирует пачками данных, структура логов адаптирована под фиксацию агрегированных результатов по каждому чанку:
| 1. Код (canonical_code) | 2. Уровень (log_level) | 3. Статусы (transport_statuses) | 4. Получатель (error_target) | 5. JSON для Фронтенда (ui_payload) | 6. Метрики для ClickHouse (observability_json) |
|---|---|---|---|---|---|
| MARKETING_BATCH_SUCCESS | INFO | {“grpc”: null, “http”: 200} | GOOGLE_FCM_API | null | {“metric”: “bulk_chunk_delivered”, “labels”: {“campaign”: “sale_2026”, “chunk_size”: 500, “success_count”: 498}} |
| MARKETING_DB_TIMEOUT | ERROR | {“grpc”: null, “http”: 504} | INTERNAL_SYSTEM | null | {“metric”: “bulk_db_timeout”, “labels”: {“replica”: “postgres_ro_pool”, “chunk_index”: 142}} |
| MARKETING_FCM_LIMIT | WARN | {“grpc”: null, “http”: 429} | GOOGLE_FCM_API | null | {“metric”: “bulk_rate_limited”, “labels”: {“campaign”: “sale_2026”, “action”: “backoff_retry”}} |