Microservice: Сервис Маркетинговых Рассылок (Marketing Campaign Worker)

Спецификация изолированного Highload-воркера пакетной дистрибуции пуш-уведомлений на Go

Author

Системный аналитик

Published

September 27, 2026

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)

  1. Чанкование вектора: Полученный из Kafka массив user_ids (размером до 1 000 000 элементов) внутри горутины нарезается на мелкие слайсы (чанки) строго по 500 штук. Это ограничение обусловлено лимитами Google API.
  2. Пакетный 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;
  3. Пакетная отправка (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”}}