Concepts

Pub/Sub

Нагрузочное тестирование pub/sub — NATS и in-memory subjects, latency публикации/доставки, проверки сообщений

Обзор

perfscale нагружает pub/sub-обмен: публикует пачки событий в subject, ждёт подходящие сообщения и измеряет как throughput публикации, так и сквозную задержку (publish → доставлено) — рядом с метриками HTTP/gRPC/WebSocket.

Pub/Sub — часть open-source движка: std/pubsub@v1 поставляется с двумя драйверами:

ДрайверТранспорт
memory (по умолчанию)Внутрипроцессная broadcast-шина — брокер не нужен. Один канал на subject, общий для всех VU в процессе: publish одного VU доходит до каждого VU, подписанного на этот subject (fan-out между VU — это фича)
natsНастоящий сервер NATS (core NATS, без JetStream)

Вендорские драйверы — Kafka, Redis, MQTT — это pro-возможность, регистрируемая через тот же driver seam (pro-драйверы). В OSS-сборке driver: kafka падает со списком зарегистрированных драйверов — всегда понятно, какая сборка нужна.

One-shot шаг

Один шаг может публиковать, потреблять или и то и другое. Когда указано оба, подписка устанавливается первой — roundtrip на том же subject всегда видит свои сообщения:

steps:
  - name: order events roundtrip
    use: std/pubsub@v1
    with:
      subject: orders.created          # memory-драйвер (по умолчанию) — без брокера
      publish:
        - '{"id":"ord-1","total":42.50}'
        - '{"id":"ord-2","total":17.00}'
      subscribe:
        count: 2                       # ждём оба сообщения
        until_contains: '"id"'         # каждое засчитанное должно совпасть
        timeout_ms: 2000
    check:
      body_contains: ord-2
ПараметрПо умолчаниюОписание
drivermemorymemory, nats или pro-драйвер (kafka, redis, mqtt)
subjectобязателенNATS subject / имя in-memory топика
urlURL брокера; обязателен для nats, игнорируется в memory
publishОдно сообщение или список; не-строки сериализуются в JSON
subscribe.count1Сколько сообщений ждать
subscribe.until_containsПодстрока, которую должно содержать каждое засчитанное сообщение
subscribe.timeout_ms5000Через столько миллисекунд шаг падает
optionsНастройки драйвера, передаются как есть (встроенные игнорируют; pro-драйверы берут отсюда QoS, consumer group, auth)

Шаг падает при ошибке подключения, ошибке публикации или таймауте подписки — в ошибке указано, сколько сообщений из count пришло и сколько отклонил matcher.

Нагрузка «только producer»

Долбим брокер без потребления — успех означает, что все публикации приняты:

steps:
  - name: produce order events
    use: std/pubsub@v1
    with:
      driver: nats
      url: nats://nats.internal:4222
      subject: orders.created
      publish:
        - '{"id":"ord-${seq}","total":${randf(10,100,2)}}'

Пayloads интерполируются, как в любом другом шаге: ${seq}, ${randf(1.05, 1.15, 5)}, outputs предыдущих шагов.

Consumer с проверками

Ждём конкретное событие и проверяем payload:

steps:
  - name: await shipment event
    use: std/pubsub@v1
    with:
      driver: nats
      url: nats://nats.internal:4222
      subject: orders.shipped
      subscribe:
        count: 1
        until_contains: '"id":"ord-1"'
        timeout_ms: 3000
    outputs: shipment
    check:
      body_contains: ord-1

outputs отдаёт published, received, duration_ms и body (matched payloads, склеенные переводом строки) для следующих шагов.

Метрики и thresholds

Каждый обмен попадает в итоговую сводку как custom-метрики: counters pubsub_msgs_published / pubsub_msgs_received и trend pubsub_e2e_ms (по сэмплу на matched-сообщение: начало фазы publish → доставлено). Гейтайте прогон через std/thresholds@v1:

  - use: std/thresholds@v1
    with:
      pubsub_e2e_ms:
        - "p(95)<50"              # 95% сообщений доставляются за 50 мс
      pubsub_e2e_ms_failed:
        - "rate<0.01"             # меньше 1% таймаутов/ошибок обмена

Под нагрузкой memory-драйвер показывает «потолок» движка (loopback, без брокера), а тот же тест против nats добавляет реальный round trip брокера — разница между ними изолирует стоимость брокера.