Sink 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.
Sinks forward processed stream data to external storage, message brokers, or network endpoints.
Development
To create a sink plugin, implement the api.Sink interface and export it from a Go plugin.
Before developing, configure the plugin development environment.
Sinks belong to two categories based on payload encoding:
BytesCollector: Receives serialized binary payloads (such as MQTT sink).TupleCollector: Receives structured map tuples and handles internal serialization (such as SQL sink).
General Methods
All sink implementations must provide these methods:
Provision:
goProvision(ctx StreamContext, configs map[string]any) errorInitializes the sink instance with configuration properties from the rule action definition (such as host, port, credentials).
Connect:
goConnect(ctx StreamContext, sch StatusChangeHandler) errorEstablishes the connection to the external destination. Reconnection logic should run asynchronously and notify connection status changes through the status handler callback.
Collect:
Receives data from upstream operators and writes it to the target system.
Close:
goClose(ctx StreamContext) errorTerminates active connections and releases resources when the rule terminates.
Export the Symbol:
Export a constructor function at the end of the file:
gofunc MySink() api.Sink { return &mySink{} }
Sink Type Implementations
BytesCollector:
goCollect(ctx StreamContext, item RawTuple) errorExtract serialized bytes with
item.Raw(). To enable automatic retry, return error messages starting with"io error".TupleCollector:
goCollect(ctx StreamContext, item MessageTuple) error CollectList(ctx StreamContext, items MessageTupleList) errorProcesses single structured records or batch lists of tuples.
Updatable Sinks
If the sink supports update or delete mutations, inspect the rowkindField property during Provision. In Collect, extract the action string (insert, update, upsert, or delete) to format the corresponding target command.
Dynamic Properties
To evaluate template expressions in sink properties at runtime, use the dynamic properties helper:
func Collect(ctx StreamContext, item RawTuple) error {
if dp, ok := item.(api.HasDynamicProps); ok {
temp, transformed := dp.DynamicProps("propName")
if transformed {
tpc = temp
}
}
return nil
}Usage in Rules
Specify the custom sink by name in the rule actions definition:
{
"id": "rule1",
"sql": "SELECT demo.temperature, demo1.temp FROM demo LEFT JOIN demo1 ON demo.timestamp = demo1.timestamp WHERE demo.temperature > demo1.temp GROUP BY demo.temperature, HOPPINGWINDOW(ss, 20, 10)",
"actions": [
{
"mySink": {
"server": "tcp://47.52.67.87:1883",
"topic": "demoSink"
}
}
]
}