Source Extension
Note
Go C-shared native .so dynamic plugins are unsupported in rekuiper. High-performance connectors are compiled directly into the rekuiper Rust binary. Custom function extensions run through WebAssembly (Wasm) or External Services. This guide is preserved as a technical reference for legacy eKuiper installations.
Sources ingest data from external systems into rekuiper streams and tables.
Sources belong to two categories:
- Scan Source: Ingests continuous data streams or scans entire tables.
- Lookup Source: Performs key-based lookups during table joins.
Develop a Scan Source
To create a scan source, implement the api.Source interface and export it from a Go plugin.
Before developing, configure the plugin development environment.
Scan sources implement one of four interface types based on ingestion model and data representation:
ByteSource: Push-based source receiving raw binary bytes. rekuiper decodes the payload according to stream format configurations.TupleSource: Push-based source that decodes custom payloads internally and emits structured map tuples.PullBytesSource: Pull-based source periodically polling external systems for binary payloads.PullTupleSource: Pull-based source periodically polling external systems and emitting structured map tuples.
General Methods
All source plugins must implement these lifecycle methods:
Provision:
goProvision(ctx StreamContext, configs map[string]any) errorInitializes the source using configuration parameters parsed from the source YAML file.
Connect:
goConnect(ctx StreamContext, sch StatusChangeHandler) errorEstablishes the connection to the external data source. Connection state changes notify the engine through the status handler callback.
Ingestion Logic:
Implement
SubscribeorPullaccording to the source interface type.Close:
goClose(ctx StreamContext) errorCloses open connections and releases resources when the rule terminates.
Export the Symbol:
Export a constructor function at the end of the file:
gofunc MySource() api.Source { return &mySource{} }
Source Interface Implementations
ByteSource:
goSubscribe(ctx StreamContext, ingest BytesIngest, ingestError ErrorIngest) errorSubscribes to external notifications and forwards raw bytes through
BytesIngest.TupleSource:
goSubscribe(ctx StreamContext, ingest TupleIngest, ingestError ErrorIngest) errorSubscribes to external notifications and forwards decoded map objects through
TupleIngest.PullBytesSource:
goPull(ctx StreamContext, trigger time.Time, ingest BytesIngest, ingestError ErrorIngest)Polls external systems at configured intervals and forwards binary payloads.
PullTupleSource:
goPull(ctx StreamContext, trigger time.Time, ingest TupleIngest, ingestError ErrorIngest)Polls external systems at configured intervals and forwards decoded map objects.
Develop a Lookup Source
A lookup source implements the api.LookupSource interface to query external systems during join operations.
LookupSource: Decodes records internally and returns maps:goLookup(ctx StreamContext, fields []string, keys []string, values []any) ([]map[string]any, error)LookupBytesSource: Returns binary payloads for automatic decoding:goLookup(ctx StreamContext, fields []string, keys []string, values []any) ([][]byte, error)
Source Traits
Extended sources can implement optional trait interfaces:
- Rewindable Source (
api.Rewindable): Required for checkpointing and exactly-once processing. ImplementGetOffset()with thread-safe synchronization. - Bounded Source (
api.Bounded): EmitsEOFIngestwhen reading completes. The engine terminates the rule automatically upon receiving the end-of-file signal.
Configuration and Usage
Source configuration files reside under etc/sources/{sourceName}.yaml.
Common parameters:
interval: Polling interval in milliseconds for pull sources.bufferLength: Maximum queue capacity in memory. Default value is102400.
To use the custom source, declare it in the stream definition:
CREATE STREAM demo (
USERID BIGINT,
FIRST_NAME STRING,
LAST_NAME STRING,
NICKNAMES ARRAY(STRING),
Gender BOOLEAN,
ADDRESS STRUCT(STREET_NAME STRING, NUMBER BIGINT)
) WITH (DATASOURCE="mytopic", TYPE="mySource", CONF_KEY="democonf");