Source Connectors Overview
Source connectors ingest data from external systems into the rekuiper stream processing engine.
Ingestion Modes
Source connectors support two ingestion modes:
- Scan Mode: Ingests event data continuously as an unbounded stream. Use this mode for stream definitions and scan tables.
- Lookup Mode: Queries external records on demand when a rule executes a join. Use this mode for lookup tables.
Connectors support one or both ingestion modes. Documentation pages indicate supported modes with badges.
Built-in Sources
rekuiper includes the following built-in source connectors:
- MQTT source: Subscribes to MQTT topics.
- RabbitMQ source: Consumes messages from RabbitMQ queues using AMQP 0-9-1.
- EdgeX source: Ingests sensor readings and events from the EdgeX message bus.
- HTTP pull source: Periodically pulls data from HTTP endpoints.
- HTTP push source: Ingests data pushed to rekuiper through HTTP POST requests.
- WebSocket source: Ingests real-time events over WebSocket connections.
- Redis source: Reads Redis keys or queries Redis as a lookup table.
- RedisSub source: Subscribes to Redis pub/sub channels.
- Kafka source: Consumes stream records from Apache Kafka topics.
- SQL source: Periodically queries relational databases through SQL.
- File source: Reads data from local files or directories.
- Memory source: Reads from internal in-memory topics to create rule pipelines.
- Simulator source: Generates synthetic sensor telemetry for testing.
NOTE
Legacy Go C-shared dynamic plugins (.so) are not supported in rekuiper. Connectors like Kafka, SQL, and WebSocket are compiled natively into the binary. For custom source extensions, use WebAssembly (Wasm) or an external HTTP service.
Using Sources in Streams and Tables
To use a source connector, create a stream or table and set the TYPE property in the WITH clause.
You can configure source behavior during stream creation by setting properties such as serialization format and decompression. For property definitions and DDL syntax, refer to Stream Management.
Runtime Execution Nodes
In a rule topology, a data source begins as a logical node. At runtime, the planner expands the logical source into an execution pipeline composed of multiple physical nodes.
Splitting source processing into distinct nodes provides three benefits:
- Component Reuse and Modularity: Shared operations (such as decompression and payload decoding) execute in standard reusable nodes.
- Sub-Task Observability: Processing stages expose fine-grained metrics for latency and throughput.
- Parallel Execution: Independent stages can execute concurrently across pipeline threads.
Source Execution Pipeline
The physical execution plan arranges source operations into sequential stages:
Connector --> RateLimit --> Decompress --> Decode --> PreprocessThe planner creates pipeline nodes based on stream configurations:
- Connector: Connects to the external system and receives raw data. Every source creates this node.
- RateLimit: Applies to push sources (such as MQTT) when the
intervalproperty is configured. This node limits the ingestion frequency. For details, refer to Down Sampling. - Decompress: Applies when the source ingests compressed binary payloads and the
decompressproperty is configured. - Decode: Deserializes raw bytes into records based on the configured
formatand schema definitions. - Preprocess: Validates records and converts data types when a stream specifies a schema and enables
strictValidation.