EdgeX Source Connector
stream sourcescan table source
The EdgeX source connector subscribes to messages from the EdgeX message bus and routes them into the rekuiper stream processing engine.
The connector processes events without manual schema definitions by using the predefined data types in EdgeX reading objects.
The EdgeX connector operates as both a source connector and a sink connector. This document explains source connector configuration and usage.
Configuration Overview
Configure the connector using environment variables, the REST API, or the configuration file.
The default configuration file is $rekuiper/etc/sources/edgex.yaml. Settings defined in the default section serve as global defaults. Custom configurations in separate sections override default values.
Example configuration file:
# Global EdgeX configurations
default:
protocol: tcp
server: localhost
port: 5573
topic: rules-events
messageType: event
# optional:
# ClientId: client1
# Username: user1
# Password: password
# Override global configurations
demo1:
protocol: tcp
server: 10.211.55.6
port: 5571
topic: rules-eventsGlobal Configurations
Properties in the default section apply to all EdgeX connections unless explicitly overridden.
Connection Parameters
protocol: Protocol used to connect to the EdgeX message bus. Default istcp.server: Server host address of the EdgeX message bus. Default islocalhost.port: Port number of the EdgeX message bus. Default is5573.
Connection Reuse
connectionSelector: Specifies a named connection profile fromconnections/connection.yaml(for example,edgex.redisMsgBus). For details, refer to Connection Management.
default:
protocol: tcp
server: localhost
port: 5573
connectionSelector: edgex.redisMsgBus
topic: rules-events
messageType: eventNOTE
When connectionSelector is specified, the engine ignores inline connection parameters (protocol, server, and port).
Topic and Message Bus Parameters
topic: EdgeX message bus topic name. Default isrules-events. SetmessageTypeto match the target topic format.type: Message bus backend type:redis: Uses Redis as the message bus. This is the default setting in EdgeX Docker Compose environments.mqtt: Uses an MQTT broker as the message bus. Configure parameters inoptional.zero: Uses ZeroMQ as the message bus.nats-jetstream: Uses NATS JetStream.nats-core: Uses NATS Core.
messageType: EdgeX payload data model:event: Decodes payloads asdtos.Event. Use this setting when subscribing to EdgeX application service topics. This is the default setting.request: Decodes payloads asrequests.AddEventRequest. Use this setting when subscribing directly to core-data or device-service buses.
Optional Parameters for MQTT Message Bus
When type is set to mqtt, configure connection settings under optional. Enclose all values in quotation marks:
ClientIdUsernamePasswordQosKeepAliveRetainedConnectionPayloadCertFileKeyFileCertPEMBlockKeyPEMBlockSkipCertVerify
Custom Configurations
Define custom configuration sections in edgex.yaml for specific topics or broker addresses:
demo1:
protocol: tcp
server: 10.211.55.6
port: 5571
topic: rules-eventsReference the configuration with CONF_KEY="demo1" in the stream DDL statement:
CREATE STREAM demo1 () WITH (FORMAT = "JSON", TYPE = "edgex", CONF_KEY = "demo1");Create a Stream Source
The EdgeX connector functions as a stream source or as a scan table source.
Create Stream via REST API
Send a POST request to /streams:
{
"sql": "CREATE STREAM demo1 () WITH (FORMAT = \"JSON\", TYPE = \"edgex\", CONF_KEY = \"demo1\")"
}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 demo '() WITH (FORMAT = "json", DATASOURCE = "demo", TYPE = "edgex")'For CLI command syntax, refer to Streams Management with CLI.
Stream Definition for EdgeX
Define EdgeX streams as schemaless streams (CREATE STREAM demo ()). EdgeX readings include predefined type information in reading objects.
Automatic Data Type Conversion
rekuiper converts reading values automatically based on the EdgeX ValueType property:
- If the engine detects a matching data type, it converts the reading value.
- If no match exists, the original value remains unchanged.
- If type conversion fails, the engine drops the value and logs a warning.
Boolean
When ValueType is Bool, rekuiper converts the value to a boolean:
- Values converted to
true:"1","t","T","true","TRUE","True" - Values converted to
false:"0","f","F","false","FALSE","False"
Bigint
When ValueType is INT8, INT16, INT32, INT64, UINT, UINT8, UINT16, UINT32, or UINT64, rekuiper converts the value to bigint.
Float
When ValueType is FLOAT32 or FLOAT64, rekuiper converts the value to float.
String
When ValueType is String, rekuiper converts the value to string.
Array Types
Boolarrays convert tobooleanarrays.INT8,INT16,INT32,INT64,UINT,UINT8,UINT16,UINT32, andUINT64arrays convert tobigintarrays.FLOAT32andFLOAT64arrays convert tofloatarrays.