Pro Features
Pub/Sub — Redis, MQTT, Kafka
Нагрузочное тестирование брокеров сообщений — драйверы Redis Pub/Sub, MQTT и Kafka для pub/sub-шага
Обзор
perfscale нагружает брокеры сообщений шагом std/pubsub@v1: публикует
пачки сообщений, читает их с проверками содержимого и измеряет сквозную
задержку (publish → consumed). Сам шаг и его драйверы memory/nats —
open 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.qos | 0 | 0, 1 или 2 — at-most-once правильный дефолт для нагрузки |
options.client_id | perfscale-<uuid> | Client ID на обмен |
options.keep_alive_secs | 30 | MQTT 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_id | perfscale-<uuid> | Клиентский group.id (клиент требует его наличия; коммитов не бывает никогда) |
options.auto_offset_reset | latest | latest — каждая партиция стартует с watermark на момент начала обмена (подходит subscribe-first roundtrip); earliest проигрывает топик с начала |
options.message_timeout_ms | 5000 | Таймаут доставки одного сообщения |
options.key | — | Ключ записи для каждого публикуемого сообщения |
options.security_protocol | — | plaintext, ssl, sasl_plaintext, sasl_ssl |
options.sasl_mechanism | — | PLAIN, SCRAM-SHA-256, SCRAM-SHA-512 |
options.sasl_username / options.sasl_password | — | SASL-креды |
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 нельзя детерминированно дождаться, а ленивый
latestreset может пропустить собственные сообщения 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-драйверы протоколов.