Execute AI Algorithms with Python Function Plugins
By integrating rekuiper and TensorFlow Lite, you can analyze streaming data by using pre-trained machine learning models. This tutorial explains how to build a Python portable plugin that classifies streaming images captured by edge devices.
You can download the completed plugin archive and source code from the eKuiper resources repository.
Prerequisites
Download a trained TensorFlow Lite model before starting. This tutorial uses the model from the TensorFlow image classification example.
Prepare the following environment:
Install Python 3.x.
Install required packages:
shellpip install pynng ekuiper tflite_runtime
By default, rekuiper starts portable plugins by using the python command. If your environment requires python3, configure the command name in the portable plugin configuration file.
When building with Docker, use the lfedge/ekuiper:<tag>-slim-python container image, which includes both rekuiper and the Python runtime.
Develop the Plugin
We will develop a function plugin named labelImage. The function takes binary image data as input and returns a string representing the recognized label. For example, when an image contains a peacock, labelImage(col) outputs peacock.
Implement the Inference Logic
- Download the Image Classification Model archive, extract its contents, and place
mobilenet_v1_1.0_224.tfliteandlabels.txtin your project folder. - Create
label.pyand implement thelabel(file_bytes)function:
import base64
import json
import tflite_runtime.interpreter as tf
def label(file_bytes):
# Load the model
interpreter = tf.Interpreter(model_path="mobilenet_v1_1.0_224.tflite")
interpreter.allocate_tensors()
input_details = interpreter.get_input_details()
output_details = interpreter.get_output_details()
# Preprocess image bytes and populate input tensors (omitted for brevity)
interpreter.set_tensor(input_details[0]['index'], input_data)
interpreter.invoke()
output_data = interpreter.get_tensor(output_details[0]['index'])
# Post-process probabilities and return label results
return resultYou can test the logic independently by adding a test script:
if __name__ == '__main__':
with open("peacock.jpg", "rb") as f:
result = label(base64.b64encode(f.read()))
print(json.dumps(result))Expected output ranked by confidence:
[
{"confidence": 0.9999935626983643, "label": "85:peacock"},
{"confidence": 2.156877371817245e-06, "label": "8:cock"},
{"confidence": 1.5930896779536852e-06, "label": "81:black grouse"}
]Implement the Plugin Interface
Create label_func.py to wrap the inference logic in the rekuiper Python plugin SDK:
from typing import List, Any
from ekuiper import Function, Context
from label import label
class LabelImageFunc(Function):
def __init__(self):
pass
def validate(self, args: List[Any]):
if len(args) != 1:
return "invalid argument length: expected 1 argument"
return ""
def exec(self, args: List[Any], ctx: Context):
return label(args[0])
def is_aggregate(self):
return False
labelIns = LabelImageFunc()Create a function metadata descriptor named functions/labelImage.json to enable user interface discovery in eKuiper manager.
Package the Plugin
Create
requirements.txtlisting all Python dependencies, and create an installation script namedinstall.sh:shell#!/bin/sh cur=$(dirname "$0") pip install -r "$cur/requirements.txt"Create an entry file named
pyai.py:pythonfrom ekuiper import PluginConfig, plugin from label_func import labelIns if __name__ == '__main__': c = PluginConfig("pyai", {}, {}, {"labelImage": lambda: labelIns}) plugin.start(c)Create the plugin metadata file named
pyai.json:json{ "version": "v1.0.0", "language": "python", "executable": "pyai.py", "sources": [], "sinks": [], "functions": [ "labelImage" ] }
Package all files into a ZIP archive with the following structure:
label.pylabel_func.pyrequirements.txtmobilenet_v1_1.0_224.tflitelabels.txtinstall.shpyai.pypyai.jsonfunctions/labelImage.json
Install the Plugin
Upload the ZIP package to the rekuiper host and install it through the REST API:
POST http://localhost:9081/plugins/portables
Content-Type: application/json
{
"name": "pyai",
"file": "file:///tmp/pyai.zip"
}Run the Plugin in Rules
Create the Stream
Define a stream that accepts binary payloads on MQTT topic tfdemo:
POST http://localhost:9081/streams
Content-Type: application/json
{
"sql": "CREATE STREAM tfdemo () WITH (DATASOURCE=\"tfdemo\", FORMAT=\"BINARY\")"
}Create the Rule
Create a rule that extracts the top classification label and publishes results to topic ekuiper/labels:
POST http://localhost:9081/rules
Content-Type: application/json
{
"id": "ruleTf",
"sql": "SELECT labelImage(self)[0]->label as label FROM tfdemo",
"actions": [
{
"mqtt": {
"server": "tcp://127.0.0.1:1883",
"sendSingle": true,
"topic": "ekuiper/labels"
}
}
]
}Publish Input Data
The following Go program publishes test images to the tfdemo topic:
package main
import (
"fmt"
"os"
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
func main() {
const TOPIC = "tfdemo"
images := []string{
"peacock.png",
"frog.jpg",
}
opts := mqtt.NewClientOptions().AddBroker("tcp://localhost:1883")
client := mqtt.NewClient(opts)
if token := client.Connect(); token.Wait() && token.Error() != nil {
panic(token.Error())
}
for _, image := range images {
fmt.Println("Publishing " + image)
payload, err := os.ReadFile(image)
if err != nil {
fmt.Println(err)
continue
}
if token := client.Publish(TOPIC, 0, false, payload); token.Wait() && token.Error() != nil {
fmt.Println(token.Error())
} else {
fmt.Println("Published " + image)
}
time.Sleep(1 * time.Second)
}
client.Disconnect(0)
}Verify Output
Subscribe to MQTT topic ekuiper/labels to verify predictions:
{"label": "85:peacock"}
{"label": "33:tailed frog, bell toad, ribbed toad, tailed toad, Ascaphus trui"}Conclusion
Python function plugins allow you to integrate machine learning inference into real-time SQL streaming rules. You can replace the demonstration model with custom models to support diverse edge AI workloads.