Python SDK for Portable Plugins
The Python SDK allows developers to build portable plugins in Python. It provides interfaces for source, sink, and function extensions, as well as runtime entry points to manage the plugin lifecycle.
Prerequisites
- Python 3.x runtime.
- Required packages: install with
pip install nng ekuiper.
By default, the engine executes Python plugins using the python command. You can specify a custom Python binary in the configuration file.
Development
Implement extension classes by subclassing the abstract base classes provided by the SDK.
Source Interface
from abc import abstractmethod
from ekuiper import Context
class Source(object):
"""Abstract base class for rekuiper source extensions."""
@abstractmethod
def configure(self, datasource: str, conf: dict):
"""Initializes configuration properties."""
pass
@abstractmethod
def open(self, ctx: Context):
"""Starts continuous ingestion and emits data or errors."""
pass
@abstractmethod
def close(self, ctx: Context):
"""Releases resources and stops ingestion."""
passSink Interface
from abc import abstractmethod
from typing import Any
from ekuiper import Context
class Sink(object):
"""Abstract base class for rekuiper sink extensions."""
@abstractmethod
def configure(self, conf: dict):
"""Initializes sink configuration properties."""
pass
@abstractmethod
def open(self, ctx: Context):
"""Establishes connections to target systems."""
pass
@abstractmethod
def collect(self, ctx: Context, data: Any):
"""Processes and forwards incoming records."""
pass
@abstractmethod
def close(self, ctx: Context):
"""Closes connections and releases resources."""
passSink Acknowledgments
When requireAck is enabled in the rule action, the sink must acknowledge each message before receiving the next record:
{
"id": "rulePort1",
"sql": "SELECT * FROM mqttStream",
"actions": [
{
"print": {
"requireAck": true
}
}
]
}The sink implementation calls ctx.ack_ok() on success or ctx.ack_error(msg) on failure:
def collect(self, ctx: Context, data: Any):
print("Received:", data)
ctx.ack_ok()Function Interface
from abc import abstractmethod
from typing import List, Any
from ekuiper import Context
class Function(object):
"""Abstract base class for rekuiper function extensions."""
@abstractmethod
def validate(self, args: List[Any]):
"""Validates arguments against expected signatures."""
pass
@abstractmethod
def exec(self, args: List[Any], ctx: Context) -> Any:
"""Computes the function output."""
pass
@abstractmethod
def is_aggregate(self) -> bool:
"""Specifies whether the function is an aggregate function."""
passMain Entry Program
Declare the plugin configuration and start the runtime:
from ekuiper import PluginConfig, plugin
if __name__ == "__main__":
c = PluginConfig(
"pysam",
{"pyjson": lambda: PyJson()},
{"print": lambda: PrintSink()},
{"revert": lambda: revertIns}
)
plugin.start(c)Refer to the Python SDK PySam Example for complete sample code.
Packaging and Deployment
Specify the main Python script file in the plugin JSON descriptor. Refer to Packaging Portable Plugins for details.
Virtual Environments (Conda)
To execute the plugin inside a Conda environment, configure the metadata descriptor:
{
"version": "v1.0.0",
"language": "python",
"executable": "pysam.py",
"virtualEnvType": "conda",
"env": "myenv",
"sources": [
"pyjson"
],
"sinks": [
"print"
],
"functions": [
"revert"
]
}