MQTT Protocol
How to use the MQTT support in Gatling to connect to a broker and perform checks against inbound messages.
Prerequisites
Gatling Enterprise MQTT SDK is not imported by default.
You have to manually add the following imports:
import io.gatling.javaapi.mqtt.*;
import static io.gatling.javaapi.mqtt.MqttDsl.*;import { mqtt } from "@gatling.io/mqtt";import io.gatling.javaapi.mqtt.MqttDsl.*import io.gatling.mqtt.Predef._MQTT protocol
Use the mqtt object in order to create a MQTT protocol.
MqttProtocolBuilder mqttProtocol = mqtt
// enable protocol version 3.1
.mqttVersion_3_1()
// enable protocol version 3.1.1 (default protocol version)
.mqttVersion_3_1_1()
// enable protocol version 5
.mqttVersion_5()
// broker address (default: localhost:1883)
.broker("hostname", 1883)
// if TLS should be enabled (default: false)
.useTls(true)
// Used to specify KeyManagerFactory for each individual virtual user. Input is the 0-based incremental id of the virtual user.
.perUserKeyManagerFactory(
userId -> {
try {
KeyStore keyStore = KeyStore.getInstance("PKCS12");
// P12 files stored under src/test/resources/keys
try (InputStream is =
getClass()
.getClassLoader()
.getResourceAsStream("keys/pk-" + userId + ".p12")) {
keyStore.load(is, null);
}
KeyManagerFactory kmf =
KeyManagerFactory.getInstance(KeyManagerFactory.getDefaultAlgorithm());
kmf.init(keyStore, null);
return kmf;
} catch (Exception e) {
throw new RuntimeException(e);
}
})
// clientIdentifier sent in the connect payload (of not set, Gatling will generate a random one)
.clientId("#{id}")
// if session should be cleaned during connect (default: true)
.cleanSession(true)
// optional credentials for connecting
.credentials("#{userName}", "#{password}")
// connections keep alive timeout
.keepAlive(30)
// use at-most-once QoS (default: true)
.qosAtMostOnce()
// use at-least-once QoS (default: false)
.qosAtLeastOnce()
// use exactly-once QoS (default: false)
.qosExactlyOnce()
// enable retain (default: false)
.retain(false)
// send last will, possibly with specific QoS and retain
.lastWill(
LastWill("#{willTopic}", StringBody("#{willMessage}"))
.qosAtLeastOnce()
.retain(true)
)
// max number of reconnects after connection crash (default: 3)
.reconnectAttemptsMax(1)
// reconnect delay after connection crash in millis (default: 100)
.reconnectDelay(1)
// reconnect delay exponential backoff (default: 1.5)
.reconnectBackoffMultiplier(1.5F)
// resend delay after send failure in millis (default: 5000)
.resendDelay(1000)
// resend delay exponential backoff (default: 1.0)
.resendBackoffMultiplier(2.0F)
// interval for timeout checker (default: 1 second)
.timeoutCheckInterval(1)
// check for pairing messages sent and messages received
.correlateBy(check)
// enable unmatched MQTT inbound messages buffering,
// with a max buffer size of 5
.unmatchedInboundMessageBufferSize(5);const mqttProtocol = mqtt
// enable protocol version 3.1
.mqttVersion_3_1()
// enable protocol version 3.1.1 (default protocol version)
.mqttVersion_3_1_1()
// enable protocol version 5
.mqttVersion_5()
// broker address (default: localhost:1883)
.broker("hostname", 1883)
// if TLS should be enabled (default: false)
.useTls(true)
// clientIdentifier sent in the connect payload (of not set, Gatling will generate a random one)
.clientId("#{id}")
// if session should be cleaned during connect (default: true)
.cleanSession(true)
// optional credentials for connecting
.credentials("#{userName}", "#{password}")
// connections keep alive timeout
.keepAlive(30)
// use at-most-once QoS (default: true)
.qosAtMostOnce()
// use at-least-once QoS (default: false)
.qosAtLeastOnce()
// use exactly-once QoS (default: false)
.qosExactlyOnce()
// enable retain (default: false)
.retain(false)
// send last will, possibly with specific QoS and retain
.lastWill(
LastWill("#{willTopic}", StringBody("#{willMessage}"))
.qosAtLeastOnce()
.retain(true)
)
// max number of reconnects after connection crash (default: 3)
.reconnectAttemptsMax(1)
// reconnect delay after connection crash in millis (default: 100)
.reconnectDelay(1)
// reconnect delay exponential backoff (default: 1.5)
.reconnectBackoffMultiplier(1.5)
// resend delay after send failure in millis (default: 5000)
.resendDelay(1000)
// resend delay exponential backoff (default: 1.0)
.resendBackoffMultiplier(2.0)
// interval for timeout checker (default: 1 second)
.timeoutCheckInterval(1)
// check for pairing messages sent and messages received
.correlateBy(check)
// enable unmatched MQTT inbound messages buffering,
// with a max buffer size of 5
.unmatchedInboundMessageBufferSize(5);val mqttProtocol = mqtt
// enable protocol version 3.1
.mqttVersion_3_1()
// enable protocol version 3.1.1 (default protocol version)
.mqttVersion_3_1_1()
// enable protocol version 5
.mqttVersion_5()
// broker address (default: localhost:1883)
.broker("hostname", 1883)
// if TLS should be enabled (default: false)
.useTls(true)
// Used to specify KeyManagerFactory for each individual virtual user. Input is the 0-based incremental id of the virtual user.
.perUserKeyManagerFactory { userId ->
val keyStore = KeyStore.getInstance("PKCS12")
// P12 files stored under src/test/resources/keys
javaClass.classLoader.getResourceAsStream("keys/pk-$userId.p12")!!.use {
keyStore.load(it, null)
}
KeyManagerFactory.getInstance(KeyManagerFactory.getDefaultAlgorithm()).apply {
init(keyStore, null)
}
}
// clientIdentifier sent in the connect payload (of not set, Gatling will generate a random one)
.clientId("#{id}")
// if session should be cleaned during connect (default: true)
.cleanSession(true)
// optional credentials for connecting
.credentials("#{userName}", "#{password}")
// connections keep alive timeout
.keepAlive(30)
// use at-most-once QoS (default: true)
.qosAtMostOnce()
// use at-least-once QoS (default: false)
.qosAtLeastOnce()
// use exactly-once QoS (default: false)
.qosExactlyOnce()
// enable retain (default: false)
.retain(false)
// send last will, possibly with specific QoS and retain
.lastWill(
LastWill("#{willTopic}", StringBody("#{willMessage}"))
.qosAtLeastOnce()
.retain(true)
)
// max number of reconnects after connection crash (default: 3)
.reconnectAttemptsMax(1)
// reconnect delay after connection crash in millis (default: 100)
.reconnectDelay(1)
// reconnect delay exponential backoff (default: 1.5)
.reconnectBackoffMultiplier(1.5f)
// resend delay after send failure in millis (default: 5000)
.resendDelay(1000)
// resend delay exponential backoff (default: 1.0)
.resendBackoffMultiplier(2.0f)
// interval for timeout checker (default: 1 second)
.timeoutCheckInterval(1)
// check for pairing messages sent and messages received
.correlateBy(check)
// enable unmatched MQTT inbound messages buffering,
// with a max buffer size of 5
.unmatchedInboundMessageBufferSize(5)val mqttProtocol = mqtt
// enable protocol version 3.1
.mqttVersion_3_1
// enable protocol version 3.1.1 (default protocol version)
.mqttVersion_3_1_1
// enable protocol version 5
.mqttVersion_5
// broker address (default: localhost:1883)
.broker("hostname", 1883)
// if TLS should be enabled (default: false)
.useTls(true)
// Used to specify KeyManagerFactory for each individual virtual user. Input is the 0-based incremental id of the virtual user.
.perUserKeyManagerFactory { userId =>
val keyStore = KeyStore.getInstance("PKCS12")
// P12 files stored under src/test/resources/keys
Using(getClass.getClassLoader.getResourceAsStream("keys/pk-" + userId + ".p12")) {
keyStore.load(_, null)
}
val kmf = KeyManagerFactory.getInstance(KeyManagerFactory.getDefaultAlgorithm)
kmf.init(keyStore, null)
kmf
}
// clientIdentifier sent in the connect payload (of not set, Gatling will generate a random one)
.clientId("#{id}")
// if session should be cleaned during connect (default: true)
.cleanSession(true)
// optional credentials for connecting
.credentials("#{userName}", "#{password}")
// connections keep alive timeout
.keepAlive(30)
// use at-most-once QoS (default: true)
.qosAtMostOnce
// use at-least-once QoS (default: false)
.qosAtLeastOnce
// use exactly-once QoS (default: false)
.qosExactlyOnce
// enable retain (default: false)
.retain(false)
// send last will, possibly with specific QoS and retain
.lastWill(
LastWill("#{willTopic}", StringBody("#{willMessage}"))
.qosAtLeastOnce
.retain(true)
)
// max number of reconnects after connection crash (default: 3)
.reconnectAttemptsMax(1)
// reconnect delay after connection crash in millis (default: 100)
.reconnectDelay(1)
// reconnect delay exponential backoff (default: 1.5)
.reconnectBackoffMultiplier(1.5F)
// resend delay after send failure in millis (default: 5000)
.resendDelay(1000)
// resend delay exponential backoff (default: 1.0)
.resendBackoffMultiplier(2.0F)
// interval for timeout checker (default: 1 second)
.timeoutCheckInterval(1)
// check for pairing messages sent and messages received
.correlateBy(check)
// enable unmatched MQTT inbound messages buffering,
// with a max buffer size of 5
.unmatchedInboundMessageBufferSize(5)Use per virtual user client certificates
With useTls(true), you can provide a client certificate for each virtual user with perUserKeyManagerFactory, either with a function that returns a javax.net.ssl.KeyManagerFactory (see the sample above) or, more simply, by storing all your keys in a single PKCS#12 file with one key entry (alias) per virtual user.
The password is optional.
// keys/multi-alias.p12 is a PKCS#12 file with one key entry per virtual user
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12", "password");
// no password
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12");// keys/multi-alias.p12 is a PKCS#12 file with one key entry per virtual user
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12", "password");
// no password
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12");// keys/multi-alias.p12 is a PKCS#12 file with one key entry per virtual user
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12", "password")
// no password
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12")// keys/multi-alias.p12 is a PKCS#12 file with one key entry per virtual user
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12", "password")
// no password
mqtt.perUserKeyManagerFactory("keys/multi-alias.p12")Note that:
- The file is located either on the classpath or as an absolute path on the filesystem.
- Key entries are assigned to virtual users in the alphabetical order of their aliases: the virtual user with id 1 gets the first alias, and so on.
- Each key entry is used by one single virtual user. If the file contains fewer key entries than virtual users, the run is stopped.
- On Gatling Enterprise with multiple load generators, aliases are split so that each load generator gets its own distinct subset, hence a key entry is never used by 2 different load generators.
See HTTP protocol for a sample script that generates such a keystore.
Request
Use the mqtt("requestName") method in order to create a MQTT request.
connect
Your virtual users first have to establish a connection.
mqtt("Connecting").connect();mqtt("Connecting").connect();mqtt("Connecting").connect()mqtt("Connecting").connectsubscribe
Use the subscribe method to subscribe to an MQTT topic:
mqtt("Subscribing")
.subscribe("#{myTopic}")
// optional, override default QoS
.qosAtMostOnce();mqtt("Subscribing").subscribe("#{myTopic}")
// optional, override default QoS
.qosAtMostOnce();mqtt("Subscribing").subscribe("#{myTopic}") // optional, override default QoS
.qosAtMostOnce()mqtt("Subscribing")
.subscribe("#{myTopic}")
// optional, override default QoS
.qosAtMostOncepublish
Use the publish method to publish a message. You can use the same Body API as for HTTP request bodies:
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"));mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"));mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"))mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"))MQTT checks
You can define blocking checks with await and non-blocking checks with expect.
Those can be set right after subscribing, or after publishing:
// subscribe and expect to receive a message within 100ms, without blocking flow
mqtt("Subscribing")
.subscribe("#{myTopic2}")
.expect(Duration.ofMillis(100));
// publish and await (block) until it receives a message withing 100ms
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(Duration.ofMillis(100));
// optionally, define in which topic the expected message will be received
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(Duration.ofMillis(100), "repub/#{myTopic}");
// optionally define check criteria to be applied on the matching received message
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(Duration.ofMillis(100))
.check(jsonPath("$.error").notExists());// subscribe and expect to receive a message within 100ms, without blocking flow
mqtt("Subscribing").subscribe("#{myTopic2}")
.expect({ amount: 100, unit: "milliseconds" });
// publish and await (block) until it receives a message withing 100ms
mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await({ amount: 100, unit: "milliseconds" });
// optionally, define in which topic the expected message will be received
mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await({ amount: 100, unit: "milliseconds" }, "repub/#{myTopic}");
// optionally define check criteria to be applied on the matching received message
mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await({ amount: 100, unit: "milliseconds" }, "repub/#{myTopic}")
.check(jsonPath("$.error").notExists());// subscribe and expect to receive a message within 100ms, without blocking flow
mqtt("Subscribing").subscribe("#{myTopic2}")
.expect(Duration.ofMillis(100))
// publish and wait (block) until it receives a message withing 100ms
mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(Duration.ofMillis(100))
// optionally, define in which topic the expected message will be received
mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(Duration.ofMillis(100), "repub/#{myTopic}")
// optionally define check criteria to be applied on the matching received message
mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(Duration.ofMillis(100))
.check(jsonPath("$.error").notExists())// subscribe and expect to receive a message within 100ms, without blocking flow
mqtt("Subscribing")
.subscribe("#{myTopic2}")
.expect(100.milliseconds)
// publish and await (block) until it receives a message withing 100ms
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(100.milliseconds)
// optionally, define in which topic the expected message will be received
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(100.milliseconds, "repub/#{myTopic}")
// optionally define check criteria to be applied on the matching received message
mqtt("Publishing")
.publish("#{myTopic}")
.message(StringBody("#{myPayload}"))
.await(100.milliseconds)
.check(jsonPath("$.error").notExists)You can optionally define in which topic the expected message will be received:
You can optionally define check criteria to be applied on the matching received message:
You can use waitForMessages and block for all pending non-blocking checks:
mqtt.waitForMessages().timeout(Duration.ofMillis(100));mqtt.waitForMessages().timeout({ amount: 100, unit: "milliseconds" });mqtt.waitForMessages().timeout(Duration.ofMillis(100))mqtt.waitForMessages.timeout(100.milliseconds)Processing unmatched messages
You can use processUnmatchedMessages to process inbound messages that haven’t been matched with a check and have been buffered.
By default, unmatched inbound messages are not buffered, you must enable this feature by setting the size of the buffer on the protocol with .unmatchedInboundMessageQueueSize(maxSize).
The buffer is reset when:
- sending an outbound message
- calling
processUnmatchedMessagesso we don’t present the same message twice
You can then pass your processing logic as a function.
The list of messages passed to this function is sorted in timestamp ascending (meaning older messages first).
It contains instances of type io.gatling.mqtt.action.MqttInboundMessage.
// store the unmatched messages in the Session
mqtt.processUnmatchedMessages("#{myTopic}", (messages, session) -> session.set("messages", messages));
// collect the last text message and store it in the Session
mqtt.processUnmatchedMessages(
"#{myTopic}",
(messages, session) ->
!messages.isEmpty()
? session.set("lastMessage", messages.get(messages.size() - 1).payloadUtf8String())
: session
);// store the unmatched messages in the Session
mqtt.processUnmatchedMessages("#{myTopic}", (messages, session) => session.set("messages", messages));
// collect the last text message and store it in the Session
mqtt.processUnmatchedMessages(
"#{myTopic}",
(messages, session) =>
messages.length > 0
? session.set("lastMessage", messages[messages.length - 1].payloadUtf8String())
: session
);// store the unmatched messages in the Session
mqtt.processUnmatchedMessages("#{myTopic}") { messages, session -> session.set("messages", messages) }
// collect the last text message and store it in the Session
mqtt.processUnmatchedMessages("#{myTopic}") { messages, session ->
messages
.map { m -> m.payloadUtf8String() }
.takeLast(1)
.fold(session) { _, lastMessage ->
session.set("lastMessage", lastMessage)
}
}// store the unmatched messages in the Session
mqtt.processUnmatchedMessages("#{myTopic}") {
(messages, session) => session.set("messages", messages)
}
// collect the last text message and store it in the Session
mqtt.processUnmatchedMessages("#{myTopic}") {
(messages, session) => {
val lastMessage = messages.reverseIterator.nextOption().map(_.payloadUtf8String)
lastMessage.fold(session)(m => session.set("lastMessage", m))
}
}MQTT configuration
MQTT support honors the ssl and netty configurations from gatling.conf.
Example
public class MqttSample extends Simulation {
MqttProtocolBuilder mqttProtocol = mqtt
.broker("localhost", 1883)
.correlateBy(jsonPath("$.correlationId"));
ScenarioBuilder scn = scenario("MQTT Test")
.feed(csv("topics-and-payloads.csv"))
.exec(mqtt("Connecting").connect())
.exec(mqtt("Subscribing").subscribe("#{myTopic}"))
.exec(mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"))
.expect(Duration.ofMillis(100)).check(jsonPath("$.error").notExists()));
{
setUp(scn.injectOpen(rampUsersPerSec(10).to(1000).during(60)))
.protocols(mqttProtocol);
}
}export default simulation((setUp) => {
const mqttProtocol = mqtt
.broker("localhost", 1883)
.correlateBy(jsonPath("$.correlationId"));
const scn = scenario("MQTT Test")
.feed(csv("topics-and-payloads.csv"))
.exec(mqtt("Connecting").connect())
.exec(mqtt("Subscribing").subscribe("#{myTopic}"))
.exec(mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"))
.expect({ amount: 100, unit: "milliseconds" })
.check(jsonPath("$.error").notExists()));
setUp(scn.injectOpen(rampUsersPerSec(10).to(1000).during(60)))
.protocols(mqttProtocol);
});class MqttSample : Simulation() {
val mqttProtocol = mqtt
.broker("localhost", 1883)
.correlateBy(jsonPath("$.correlationId"))
val scn = scenario("MQTT Test")
.feed(csv("topics-and-payloads.csv"))
.exec(mqtt("Connecting").connect())
.exec(mqtt("Subscribing").subscribe("#{myTopic}"))
.exec(mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"))
.expect(Duration.ofMillis(100)).check(jsonPath("$.error").notExists()))
init {
setUp(scn.injectOpen(rampUsersPerSec(10.0).to(1000.0).during(60)))
.protocols(mqttProtocol)
}
}class MqttSample extends Simulation {
val mqttProtocol = mqtt
.broker("localhost", 1883)
.correlateBy(jsonPath("$.correlationId"))
val scn = scenario("MQTT Test")
.feed(csv("topics-and-payloads.csv"))
.exec(mqtt("Connecting").connect)
.exec(mqtt("Subscribing").subscribe("#{myTopic}"))
.exec(mqtt("Publishing").publish("#{myTopic}")
.message(StringBody("#{myTextPayload}"))
.expect(100.milliseconds).check(jsonPath("$.error").notExists))
setUp(scn.inject(rampUsersPerSec(10) to 1000 during (60)))
.protocols(mqttProtocol)
}