Serialization
rekuiper uses an internal map-based data structure during stream computation. Source and sink connectors communicating with external systems require codecs to convert data formats. Specify the encoding and decoding configuration by setting format and schemaId in source or sink parameters.
Formats
rekuiper supports schema-based and schema-less serialization formats: json, binary, delimited, protobuf, and custom.
protobuf is a schema-based format. You must register the schema before referencing it in a rule.
The following configuration specifies Protobuf serialization for an MQTT sink:
{
"mqtt": {
"server": "tcp://127.0.0.1:1883",
"topic": "sample",
"format": "protobuf",
"schemaId": "proto1.Book"
}
}rekuiper supports three codec implementations:
- Built-in Codecs: Executed internally without external dependencies (for example, JSON parsing).
- Dynamic Schema Codecs: Parse schema files at runtime (for example, Protobuf reading
*.protofiles). - Static Plugin Codecs: Use compiled shared libraries (
*.so) for maximum parsing performance.
The following table summarizes supported formats and their capabilities:
| Format | Codec | Custom Codec | Schema |
|---|---|---|---|
json | Built-in | Unsupported | Unsupported |
binary | Built-in | Unsupported | Unsupported |
delimited | Built-in (specify delimiter) | Unsupported | Unsupported |
protobuf | Built-in | Supported | Supported and required |
parquet | Built-in (Arrow/Snappy) | Unsupported | Inferred automatically |
custom | Not built-in | Supported and required | Supported and optional |
Format Extensions
You can implement custom codecs and schemas for custom and protobuf formats by creating Go plugins:
Implement the
Converterinterface. TheEncodemethod serializes data into a byte array for sinks. TheDecodemethod deserializes bytes into map structures for sources:go// Converter converts bytes & map or []map according to the schema type Converter interface { Encode(d interface{}) ([]byte, error) Decode(b []byte) (interface{}, error) }Implement the
SchemaProviderinterface if the format is strongly typed. The method returns a JSON-schema representation used for SQL validation and optimization:gotype SchemaProvider interface { GetSchemaJson() string }Compile the code into a shared object plugin:
shellgo build -trimpath --buildmode=plugin -o data/test/myFormat.so internal/converter/custom/test/*.goRegister the schema by using the REST API:
httpPOST http://localhost:9081/schemas/custom Content-Type: application/json { "name": "custom1", "soFile": "file:///tmp/custom1.so" }Reference the format in sources or sinks by setting
format="custom"andschemaId="custom1".
Refer to myFormat.go for a complete sample implementation.
Build Format Plugins with Docker
Compile format plugins in an environment matching the target rekuiper binary. Official release images use Debian or Alpine Linux.
- Debian: Use the corresponding developer image (for example,
1.8.0-dev). - Alpine: Use the official Go Alpine image matching the engine version:
Create a
Makefilein your plugin repository. Refer to the sample project.Check the
GO_VERSIONargument in the Docker build file (for example,1.25.4).Compile the plugin inside the Alpine container:
shellcd ${yourProjectLoc} docker run --rm -it -v "$PWD":/usr/src/myapp -w /usr/src/myapp golang:1.25.4-alpine sh # Inside the container: apk add gcc make libc-dev makeLocate the compiled
.sofile and register it through the schema registry API.
Static Protobuf
For high-throughput requirements, compile static Protobuf plugins instead of using dynamic schema parsing:
Generate Go code from your
.protodefinition by usingprotoc:shellprotoc --go_opt=Mhelloworld.proto=com.main --go_out=. helloworld.protoMove the generated
helloworld.pb.gofile into your plugin project and set package name tomain.Create a wrapper struct for each message type. Implement
Encode,Decode, and accessor methods without reflection.Compile the plugin:
shellgo build -trimpath --buildmode=plugin -o data/test/helloworld.so internal/converter/protobuf/test/*.goRegister the schema by providing both the
.protodefinition and the.sobinary:httpPOST http://localhost:9081/schemas/protobuf Content-Type: application/json { "name": "helloworld", "file": "file:///tmp/helloworld.proto", "soFile": "file:///tmp/helloworld.so" }Reference the registered schema in stream and action definitions.
Refer to the helloworld protobuf sample for a full implementation.
Parquet Format
parquet is an open-source columnar storage format. rekuiper includes built-in support for Apache Parquet through Apache Arrow.
Parquet provides high compression ratios and fast analytical scan performance on edge devices.
Capabilities
- Columnar Storage: Data is organized into row groups and columns.
- Snappy Compression:
rekuipercompresses Parquet files with Snappy by default. - Automatic Schema Inference: The engine creates the Arrow schema automatically from your stream data. You do not need to register a schema file.
File Sink Example
To write stream outputs to a Parquet file, set format to "parquet" in the file action:
{
"file": {
"path": "/kuiper/data/telemetry.parquet",
"format": "parquet",
"rollingCount": 1000
}
}File Source Example
To read records from a Parquet file, create a stream with FORMAT = "PARQUET":
CREATE STREAM parquet_history () WITH (
TYPE = "file",
FILETYPE = "parquet",
PATH = "/kuiper/data/telemetry.parquet",
FORMAT = "PARQUET"
);Schema Registry
Schemas define structured record formats. rekuiper stores schema files in data/schemas/${type} (for example, data/schemas/protobuf).
During startup, rekuiper scans the schema directory and registers all definitions automatically. Manage schemas at runtime through the Schema Registry API: