Connectors
Connectors connect the rekuiper stream processing engine to external systems, such as message brokers, databases, and file systems.
Connectors ingest data from external systems into rekuiper and dispatch processed results to target endpoints. Connectors support edge environments, local networks, and cloud infrastructures.
Connectors belong to two categories:
- Source Connectors: Ingest data from external systems into the rekuiper processing pipeline.
- Sink Connectors: Dispatch processed query results from rekuiper to external destinations.
rekuiper provides built-in connectors for common protocols and supports plugin connectors for custom protocols.
Source Connectors
Source connectors ingest external data into streams or tables. Sources operate in two modes:
- Streaming mode: Ingests sequential, unbounded events in real time.
- Reference mode: Fetches batch snapshots or performs key lookups (used with tables).
To configure a source, specify the source type in the WITH clause of the stream or table definition.
Built-in Source Connectors
rekuiper includes the following built-in source connectors:
- MQTT source: Ingests messages from MQTT topics.
- EdgeX source: Ingests event data from EdgeX Foundry message buses.
- HTTP pull source: Pulls data periodically from HTTP endpoints.
- HTTP push source: Receives incoming HTTP requests sent to rekuiper endpoints.
- File source: Reads data from local files. Commonly used as tables.
- Memory source: Consumes events from internal memory topics to create rule pipelines.
- Redis source: Queries key-value data in Redis as a lookup table.
Plugin Source Connectors
Use plugin connectors for specialized data protocols or custom systems:
- SQL source: Queries relational databases on a schedule or as a lookup table.
- Video source: Captures image frames from video streams.
- Random source: Generates synthetic mock data for functional testing.
- ZeroMQ source: Consumes messages from ZeroMQ publishers.
- Kafka source: Consumes event streams from Apache Kafka topics.
Sink Connectors
Sink connectors transfer processed output records to external endpoints. Sinks support disk caching to manage network interruptions and prevent data loss. Sinks also support dynamic properties and shared connection pools.
Built-in Sink Connectors
rekuiper includes the following built-in sink connectors:
- MQTT sink: Publishes messages to an external MQTT broker.
- EdgeX sink: Sends events to EdgeX Foundry. Available when compiled with the EdgeX build tag.
- REST sink: Sends HTTP requests to external web servers.
- Redis sink: Writes key-value records and data structures to Redis.
- File sink: Writes output records to local files.
- Memory sink: Publishes records to internal memory topics to feed downstream rules.
- Log sink: Writes output records to system log files for diagnostic debugging.
- Nop sink: Discards output records without I/O operations for performance benchmarking.
Plugin Sink Connectors
Use plugin sink connectors for external platforms and custom targets:
- InfluxDB sink: Writes time-series points to InfluxDB v1.x.
- InfluxDB v2 sink: Writes time-series points to InfluxDB v2.x.
- Image sink: Writes binary image frames to local storage.
- ZeroMQ sink: Publishes messages to ZeroMQ subscribers.
- Kafka sink: Produces messages to Apache Kafka topics.
Data Templates in Sinks
Data templates transform output payloads to match target external formats. Data templates use the Golang text template syntax to support field mapping, conditional formatting, and iteration.
Batch Configuration
rekuiper supports batch import and export of connector, stream, and rule configurations through the REST API:
{
"streams": {},
"tables": {},
"rules": {},
"nativePlugins": {},
"portablePlugins": {},
"sourceConfig": {},
"sinkConfig": {}
}For configuration import and export procedures, refer to Data Import and Export Management.