Pro Features

Pub/Sub — Redis, MQTT, Kafka

Нагрузочное тестирование брокеров сообщений — драйверы Redis Pub/Sub, MQTT и Kafka для pub/sub-шага

Обзор

perfscale нагружает брокеры сообщений шагом std/pubsub@v1: публикует пачки сообщений, читает их с проверками содержимого и измеряет сквозную задержку (publish → consumed). Сам шаг и его драйверы memory/natsopen source; драйверы на этой странице — Redis Pub/Sub, MQTT и Kafkaплатная возможность (тарифы Scale и Enterprise), которую платные агенты регистрируют на том же драйверном шве.

В OSS-сборке driver: kafka падает со списком зарегистрированных драйверов — всегда понятно, какая сборка нужна.

Все драйверы разделяют семантику шага:

  • subscribe-first — если заданы и publish, и subscribe, подписка устанавливается до публикации, поэтому roundtrip на одном subject всегда видит собственные сообщения;
  • требуется хотя бы одно из publish / subscribe — publish-only это чистый producer-шаг, subscribe-only — чистый consumer;
  • таймаут подписки роняет шаг с диагностикой, сколько сообщений из count пришло и сколько отбросил матчер until_contains — это провал проверки, а не транспортная ошибка;
  • счётчики pubsub_msgs_published / pubsub_msgs_received и тренд pubsub_e2e_ms попадают в сводку прогона и в thresholds.

Настройки конкретного драйвера передаются в объекте options шага и уходят драйверу как есть. Неизвестные ключи отклоняются со списком поддерживаемых — опечатка падает явно, а не игнорируется молча. Числовые опции принимают и строковую форму (qos: "1"), поэтому работают значения из ${{ }}-интерполяции.

driver: redis

Redis Pub/Sub. subject — это Pub/Sub-канал; аутентификация живёт в URL, поэтому опций у драйвера нет.

ПараметрПо умолчаниюОписание
urlобязателенredis://[:password@]host:6379[/db] или rediss:// для TLS

Roundtrip с проверками содержимого:

steps:
  - name: order events roundtrip
    use: std/pubsub@v1
    with:
      driver: redis
      url: redis://:secret@redis.internal:6379/0
      subject: orders.created
      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

Redis Pub/Sub — fire-and-forget: ни персистентности, ни replay — подписчик, который ещё не подключился, пропускает сообщения. Шаг подписывается первым, поэтому roundtrip внутри шага безопасен; для consumer-only шага убедитесь, что producer начинает публиковать после того, как поднялись consumer'ы.

Пара producer / consumer

Реалистичный паттерн для брокера: один процесс генерирует нагрузку, другие читают. Два файла, один общий конфиг, запуск параллельно на один брокер:

# producer.yaml — чистый producer, успех = каждая публикация принята
steps:
  - name: produce order events
    use: std/pubsub@v1
    with:
      driver: redis
      url: redis://127.0.0.1:6379
      subject: orders.created
      publish:
        - '{"id":"ord-${seq}","total":${randf(10,100,2)}}'
# consumer.yaml — чистый consumer: ждёт 10 событий за итерацию
steps:
  - name: consume order events
    use: std/pubsub@v1
    with:
      driver: redis
      url: redis://127.0.0.1:6379
      subject: orders.created
      subscribe:
        count: 10
        until_contains: '"id"'
        timeout_ms: 30000
# config.yaml
vus: 4
duration: 30s
perfscale run -f consumer.yaml -c config.yaml &
perfscale run -f producer.yaml -c config.yaml

Consumer, который не успевает, роняет шаг по таймауту подписки — это и есть сигнал, под который вы сайзите консьюмеров.

driver: mqtt

MQTT v3.1.1 (самая поддерживаемая версия — v5 сознательно не используется). subject — это MQTT-топик; wildcards (sensors/+) принимаются для подписки, но норма — точные топики.

ПараметрПо умолчаниюОписание
urlобязателенmqtt://[user:pass@]host[:port] (порт по умолчанию 1883) или mqtts://… (TLS, по умолчанию 8883)
options.qos00, 1 или 2 — at-most-once правильный дефолт для нагрузки
options.client_idperfscale-<uuid>Client ID на обмен
options.keep_alive_secs30MQTT keep-alive

Roundtrip на QoS 1 с кредами:

steps:
  - name: telemetry roundtrip
    use: std/pubsub@v1
    with:
      driver: mqtt
      url: mqtt://device-sim:s3cret@mqtt.internal:1883
      subject: sensors/temperature
      publish:
        - '{"sensor":"t-1","celsius":${randf(18,26,1)}}'
      subscribe:
        count: 1
        until_contains: '"sensor"'
        timeout_ms: 3000
      options:
        qos: 1

Consumer-only через TLS, ожидание события от флота:

steps:
  - name: await overheat alert
    use: std/pubsub@v1
    with:
      driver: mqtt
      url: mqtts://mqtt.internal:8883
      subject: sensors/alerts
      subscribe:
        count: 1
        until_contains: overheat
        timeout_ms: 60000
      options:
        qos: 1
        client_id: perfscale-listener
    check:
      body_contains: overheat

mqtts:// использует rustls с webpki roots — сертификат брокера должен сходиться к публичному CA и совпадать с хостом в URL. Публикации QoS 1/2 считаются принятыми в момент передачи в event loop клиента, а не по PUBACK/PUBCOMP.

driver: kafka

Apache Kafka через rdkafka (встроенный librdkafka). subject — это топик.

ПараметрПо умолчаниюОписание
urlобязателенBootstrap-серверы: host:9092 или h1:9092,h2:9092; опциональный префикс kafka:// отрезается
options.group_idperfscale-<uuid>Клиентский group.id (клиент требует его наличия; коммитов не бывает никогда)
options.auto_offset_resetlatestlatest — каждая партиция стартует с watermark на момент начала обмена (подходит subscribe-first roundtrip); earliest проигрывает топик с начала
options.message_timeout_ms5000Таймаут доставки одного сообщения
options.keyКлюч записи для каждого публикуемого сообщения
options.security_protocolplaintext, ssl, sasl_plaintext, sasl_ssl
options.sasl_mechanismPLAIN, SCRAM-SHA-256, SCRAM-SHA-512
options.sasl_username / options.sasl_passwordSASL-креды

Roundtrip на одном топике:

steps:
  - name: order events roundtrip
    use: std/pubsub@v1
    with:
      driver: kafka
      url: kafka://kafka-1.internal:9092,kafka-2.internal:9092
      subject: orders.created
      publish:
        - '{"id":"ord-1","total":42.50}'
        - '{"id":"ord-2","total":17.00}'
      subscribe:
        count: 2
        until_contains: '"id"'
        timeout_ms: 10000

Managed-кластеры (SASL/SSL, SCRAM) — publish-only producer:

steps:
  - name: produce order events
    use: std/pubsub@v1
    with:
      driver: kafka
      url: kafka://broker-1.example.com:9092,broker-2.example.com:9092
      subject: orders.created
      publish:
        - '{"id":"ord-${seq}","total":${randf(10,100,2)}}'
      options:
        security_protocol: sasl_ssl
        sasl_mechanism: SCRAM-SHA-512
        sasl_username: load-generator
        sasl_password: ${env(KAFKA_PASSWORD)}
        message_timeout_ms: 10000

Проигрывание топика для замера пропускной способности на чтение (earliest + длинное ожидание):

steps:
  - name: drain the topic
    use: std/pubsub@v1
    with:
      driver: kafka
      url: kafka://kafka-1.internal:9092
      subject: orders.created
      subscribe:
        count: 1000
        timeout_ms: 60000
      options:
        auto_offset_reset: earliest

Две специфики Kafka, которые стоит знать:

  • Ручное назначение партиций, а не consumer-group join. До публикации стартовый offset каждой партиции фиксируется на watermark на момент начала обмена. Group join нельзя детерминированно дождаться, а ленивый latest reset может пропустить собственные сообщения roundtrip'а; зафиксированные watermark'и делают правило subscribe-first точным. Каждый обмен видит все партиции, поэтому consumer'ы фан-аутятся как в NATS-драйвере, а enable.auto.commit всегда выключен — нагрузочный прогон никогда не двигает offset'ы реальной consumer-группы.
  • Топики создаются best-effort (replication factor 1) перед чтением, чтобы у roundtrip'а на свежем топике было что назначать. Ошибки по топику (уже существует, нет прав) игнорируются; кластеры с автосозданием или заранее созданными топиками ведут себя так же.

Thresholds

Прогоны по брокерам гейтятся теми же кастомными метриками, что и у OSS-драйверов — перцентили e2e-задержки и доля ошибок:

  - use: std/thresholds@v1
    with:
      pubsub_e2e_ms:
        - "p(95)<100"             # 95% сообщений доезжают за 100 мс
      pubsub_e2e_ms_failed:
        - "rate<0.01"             # меньше 1% обменов с таймаутом/ошибкой
      pubsub_msgs_received:
        - "count>=10000"          # consumer'ы реально успевают

Модель подключений

Как и встроенный nats-драйвер, каждый обмен — это полный цикл: connect → subscribe → publish → collect → drop. Подключения не пулятся между шагами и итерациями VU — под нагрузкой каждый VU открывает свежее соединение с брокером на каждую итерацию. Это честно для e2e-замера (стоимость connect включена), но брокер видит реальный connection churn; сайзьте частоту запросов соответственно и предпочитайте прогоны, ограниченные duration, когда целитесь в общий staging-брокер.

См. также

  • Концепция Pub/Sub — параметры шага, вывод и open-source драйверы memory / nats.
  • SOAP и FIX — другие pro-драйверы протоколов.