MQTT Source Connector
stream sourcescan table source
The MQTT source connector subscribes to messages from an MQTT broker and channels them into the rekuiper stream processing pipeline.
The MQTT connector functions as both a source connector and a sink connector. This document describes source connector configuration and usage.
Configuration Overview
Configure the MQTT connector through environment variables, the REST API, or the configuration file.
The default configuration file resides at $rekuiper/etc/mqtt_source.yaml. Settings defined in the default section serve as global defaults. Custom configurations in separate sections override default values.
Example configuration file with default and demo_conf sections:
# Global MQTT configurations
default:
qos: 1
server: "tcp://127.0.0.1:1883"
#username: user1
#password: password
#certificationPath: /var/kuiper/xyz-certificate.pem
#privateKeyPath: /var/kuiper/xyz-private.pem.key
#rootCaPath: /var/kuiper/xyz-rootca.pem
#insecureSkipVerify: true
#connectionSelector: mqtt.mqtt_conf1
#decompression: ""
# Override global configurations
demo_conf:
qos: 0
server: "tcp://10.211.55.6:1883"Global Configurations
Properties in the default section apply to all MQTT connections unless explicitly overridden.
Connection Parameters
qos: Default subscription QoS level (0,1, or2). Default is1.server: Target MQTT broker URL.username: Username credential for broker authentication.password: Password credential for broker authentication.protocolVersion: MQTT protocol version:3.1(MQTT 3),3.1.1(MQTT 4), or5(MQTT 5). Default is3.1.clientid: Client identifier for the connection. If omitted, the engine generates a random UUID.
Security and TLS Parameters
certificationPath: Path to the client certificate file (for example,d3807d9fa5-certificate.pem). Can be absolute or relative to the execution root directory.privateKeyPath: Path to the client private key file (for example,d3807d9fa5-private.pem.key).rootCaPath: Path to the Root CA certificate file.certficationRaw: Base64-encoded client certificate text. The engine preferscertificationPathif both are defined.privateKeyRaw: Base64-encoded client private key text. The engine prefersprivateKeyPathif both are defined.rootCARaw: Base64-encoded Root CA certificate text. The engine prefersrootCaPathif both are defined.insecureSkipVerify: Boolean. Set totrueto skip certificate and hostname validation.
For mTLS procedures and secret handling, refer to the Secure MQTT with TLS Guide.
Connection Reuse
connectionSelector: Specifies a named connection resource fromconnections/connection.yaml(for example,mqtt.localConnection). For details, refer to Connection Management.
default:
qos: 1
server: "tcp://127.0.0.1:1883"
connectionSelector: mqtt.localConnectionNOTE
When connectionSelector is configured in a configuration group, the engine ignores broker connection parameters (such as server) defined in that group.
Verify broker reachability before runtime using the Connectivity Check API.
Payload Handling
decompression: Decompresses incoming binary payloads before parsing. Supported algorithms:"gzip"and"zstd".bufferLength: Maximum number of messages buffered in memory to prevent out-of-memory errors. Default is102400.
KubeEdge Integration
kubeedgeVersion: KubeEdge version number.kubeedgeModelFile: KubeEdge model template filename located inetc/sources/.
Example model file:
{
"deviceModels": [{
"name": "device1",
"properties": [{
"name": "temperature",
"dataType": "int"
}, {
"name": "temperature-enable",
"dataType": "string"
}]
}]
}deviceModels.name: Device name matched against the third and fourth segments of topic$ke/events/device/device1/data/update.properties.name: Property field name.properties.dataType: Expected property data type.
Custom Configurations
Define custom configuration sections in etc/mqtt_source.yaml to specify connection settings for distinct brokers or topics:
demo_conf:
qos: 0
server: "tcp://10.211.55.6:1883"To apply this configuration, specify CONF_KEY="demo_conf" in the stream definition:
CREATE STREAM demo () WITH (DATASOURCE="test/", FORMAT="JSON", KEY="USERID", CONF_KEY="demo_conf");Properties in demo_conf override corresponding values in the default section.
Create a Stream Source
The MQTT connector operates as a stream source or as a scan table source.
Create Stream via REST API
Send a POST request to /streams:
{
"sql": "CREATE STREAM my_stream (id bigint, name string, score float) WITH (DATASOURCE = \"topic/temperature\", FORMAT = \"json\", KEY = \"id\")"
}For REST API specifications, refer to Streams Management with REST API.
Create Stream via CLI
Run the kuiper create stream command:
bin/kuiper create stream my_stream '(id bigint, name string, score float) WITH (DATASOURCE = "topic/temperature", FORMAT = "json", KEY = "id")'For CLI command syntax, refer to Streams Management with CLI.
MQTT v5 User Properties
When protocolVersion is set to 5, rekuiper exposes incoming MQTT v5 User Properties in record metadata under the properties key as a map of strings (map[string]string).
Access user properties in SQL queries using the meta function:
SELECT meta(properties) AS props FROM demo;Migration Notes
Starting with version 1.5.0, the MQTT source configuration parameter changed from servers (array) to server (single URL string).
- When upgrading to 1.5.0 or later, verify that
serveris configured inetc/mqtt_source.yaml. - When using environment variable overrides, replace
MQTT_SOURCE__DEFAULT__SERVERS=[tcp://127.0.0.1:1883]withMQTT_SOURCE__DEFAULT__SERVER="tcp://127.0.0.1:1883".
Listen to Multiple Topics
To subscribe to multiple MQTT topics within a single stream, specify a comma-separated list in the DATASOURCE property:
CREATE STREAM my_stream (id bigint, name string, score float)
WITH (DATASOURCE = "t1,t2", FORMAT = "json", KEY = "id");