Skip to main content

Plugin configuration demo

Last updated 10/03/2026

LogBus configuration​

Configuration explained:

  • cmd: The command that starts the plugin server
  • depend_sh: Whether the start command depends on a bash environment. Default: false (required for Python and Java)
  • hand_shake: Plugin verification
  • hand_shake.protocol_version: Must match the definition in the plugin server
  • hand_shake.magic_cookie_key: Must match the definition in the plugin server
  • hand_shake.magic_cookie_value: Must match the definition in the plugin server
{
"datasource": [
{
"type": "kafka",
"topic": "your_topic",
"consumer_group": "YOUR_CONSUMER_GROUP",
"auto_commit": false,
"brokers": [
"YOUR_KAFKA_HOST:9092"
],
"block_partitions_timeout": 4000,
"block_partitions_revoked": true,
"app_id": "YOUR_APP_ID",
"parser": {
"cmd": "./parser",
"hand_shake": {
"protocol_version": 1,
"magic_cookie_key": "LogBus",
"magic_cookie_value": "v2"
}
}
}
],
"push_url": "http://RECEIVER_URL"
}

Plugin development​

Format of data received by the gRPC plugin server​

Data is processed in batches. Multiple records are separated by double colons "::". Each record is prefixed with the file name or the Kafka topic name, joined with the characters "|ta|".

example:

File name: event.log

Two raw records: {"user_id": "abc", "distinct_id": "123"} {"user_id": "def", "distinct_id": "456"}

Data received by the gRPC plugin server: event.log|ta|{"user_id": "abc", "distinct_id": "123"} ::event.log|ta|{"user_id": "def", "distinct_id": "456"}

Format of data returned by the gRPC plugin server​

Data converted to the standard TA format, joined with double colons "::"

example:

The returned data should be:

{"#account_id": "abc", "#distinct_id": "123"}::{"#account_id": "def", "#distinct_id": "456"}

Error handling​

  1. If invalid JSON data is returned, the main program skips the erroneous rows
  2. If data that is not in the standard TE format is returned, data reporting fails
  3. If the second return value is an error, the request is retried indefinitely to prevent erroneous data from entering the TE cluster. You can use this together with the alert feature to learn about errors in advance.

Demo program (Golang)​

For the protobuf definitions, see the appendix and choose the one for your language

package main

import (
"fmt"
"github.com/hashicorp/go-plugin"
proto "github.com/hashicorp/go-plugin/examples/logbus/prot"
"golang.org/x/net/context"
"google.golang.org/grpc"
)

const separator = "::"

// Parser custom parser
type Parser interface {
Parse([]byte) ([]byte, error)
}

// GRPCServer plugin server
type GRPCServer struct {
Impl Parser
}

type GRPCPlugin struct {
plugin.Plugin
Impl Parser
}

// GRPCClient GRPC client
func (_ *GRPCPlugin) GRPCClient(_ context.Context, _ *plugin.GRPCBroker, c *grpc.ClientConn) (interface{}, error) {
return nil, nil
}

func (g *GRPCPlugin) GRPCServer(_ *plugin.GRPCBroker, s *grpc.Server) error {
proto.RegisterParserServer(s, &GRPCServer{Impl: g.Impl})
return nil
}

// Parse server-side parsing
func (g *GRPCServer) Parse(_ context.Context, req *proto.Request) (*proto.Response, error) {
data, err := g.Impl.Parse(req.Content)
if err != nil {
fmt.Println(err)
}
return &proto.Response{
Content: data,
}, nil
}

type Par struct{}

// Parse server-side parsing
func (g *Par) Parse(b []byte) ([]byte, error) {
return b, nil
}

var handshakeConfig = plugin.HandshakeConfig{
ProtocolVersion: 1,
MagicCookieKey: "LogBus",
MagicCookieValue: "v2",
}

func main() {
plugin.Serve(&plugin.ServeConfig{
HandshakeConfig: handshakeConfig,
Plugins: map[string]plugin.Plugin{
"parser": &GRPCPlugin{Impl: &Par{}},
},
GRPCServer: plugin.DefaultGRPCServer,
})
}

Demo program (Python)​

Verification

  • hand_shake.protocol_version: Required. The value is 1
  • hand_shake.magic_cookie_key: Can be ignored
  • hand_shake.magic_cookie_value: Can be ignored

print("1|1|tcp|127.0.0.1:1234|grpc")

Note this line: it writes to standard output in the required format

1: Fixed, defined in the plugin

1 protocol_version

tcp: Protocol

address:port: Address

grpc: Protocol

from concurrent import futures
import sys
import time

import grpc

import parser_pb2
import parser_pb2_grpc

from grpc_health.v1.health import HealthServicer
from grpc_health.v1 import health_pb2, health_pb2_grpc

class ParserServicer(parser_pb2_grpc.ParserServicer):
"""Implementation of Parser service."""

def Parse(self, request, context):
data = request.content
result = parser_pb2.Response()
result.content = data
return result

def serve():
# We need to build a health service to work with go-plugin
health = HealthServicer()
health.set("plugin", health_pb2.HealthCheckResponse.ServingStatus.Value('SERVING'))

# Start the server.
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
parser_pb2_grpc.add_ParserServicer_to_server(ParserServicer(), server)
health_pb2_grpc.add_HealthServicer_to_server(health, server)
server.add_insecure_port('127.0.0.1:1234')
server.start()

# Output information
print("1|1|tcp|127.0.0.1:1234|grpc")
sys.stdout.flush()

try:
while True:
time.sleep(60 * 60 * 24)
except KeyboardInterrupt:
server.stop(0)

if __name__ == '__main__':
serve()

Demo program (Java)​

package cn.thinkingdata.logbus.grpc.demo.service;

import cn.thinkingdata.logbus.grpc.demo.proto.ParserGrpc;
import cn.thinkingdata.logbus.grpc.demo.proto.ParserOuterClass;
import com.google.protobuf.ByteString;
import io.grpc.stub.StreamObserver;

import java.util.ArrayList;
import java.util.List;


public class TransferGrpcService extends ParserGrpc.ParserImplBase {

@Override
public void parse(ParserOuterClass.Request request, StreamObserver<ParserOuterClass.Response> responseObserver) {
String content = request.getContent().toStringUtf8();
String[] datas = content.split("::");
List<String> responseList = new ArrayList<>();
for (String data : datas) {
try {
String[] split = data.split("\\|ta\\|");
responseList.add(split[1]);
} catch (Exception e) {
e.printStackTrace();
}
}
String join = String.join("::", responseList);
ParserOuterClass.Response response = ParserOuterClass.Response.newBuilder().setContent(ByteString.copyFrom(join.getBytes())).build();
responseObserver.onNext(response);
responseObserver.onCompleted();
}
}

Appendix​

protobuf​

Protobuf structure files​
📎 parser.proto(1 KB)
Protobuf files for Python​
python -m grpc_tools.protoc -I ./proto/ --python_out=./plugin-python/ --grpc_python_out=./plugin-python/ ./proto/parser.proto
parser_pb2.py​
# -*- coding: utf-8 -*-
# Generated by the protocol buffer compiler. DO NOT EDIT!
# source: parser.proto
"""Generated protocol buffer code."""
from google.protobuf.internal import builder as _builder
from google.protobuf import descriptor as _descriptor
from google.protobuf import descriptor_pool as _descriptor_pool
from google.protobuf import symbol_database as _symbol_database
# @@protoc_insertion_point(imports)

_sym_db = _symbol_database.Default()




DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x0cparser.proto\x12\x05proto\"\x1a\n\x07Request\x12\x0f\n\x07\x63ontent\x18\x01 \x01(\x0c\"\x1b\n\x08Response\x12\x0f\n\x07\x63ontent\x18\x01 \x01(\x0c\x32\x32\n\x06Parser\x12(\n\x05Parse\x12\x0e.proto.Request\x1a\x0f.proto.ResponseB\nZ\x08../protob\x06proto3')

_builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, globals())
_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'parser_pb2', globals())
if _descriptor._USE_C_DESCRIPTORS == False:

DESCRIPTOR._options = None
DESCRIPTOR._serialized_options = b'Z\010../proto'
_REQUEST._serialized_start=23
_REQUEST._serialized_end=49
_RESPONSE._serialized_start=51
_RESPONSE._serialized_end=78
_PARSER._serialized_start=80
_PARSER._serialized_end=130
# @@protoc_insertion_point(module_scope)
parser_pb2_grpc.py​
# Generated by the gRPC Python protocol compiler plugin. DO NOT EDIT!
"""Client and server classes corresponding to protobuf-defined services."""
import grpc

import parser_pb2 as parser__pb2


class ParserStub(object):
"""Missing associated documentation comment in .proto file."""

def __init__(self, channel):
"""Constructor.

Args:
channel: A grpc.Channel.
"""
self.Parse = channel.unary_unary(
'/proto.Parser/Parse',
request_serializer=parser__pb2.Request.SerializeToString,
response_deserializer=parser__pb2.Response.FromString,
)


class ParserServicer(object):
"""Missing associated documentation comment in .proto file."""

def Parse(self, request, context):
"""Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')


def add_ParserServicer_to_server(servicer, server):
rpc_method_handlers = {
'Parse': grpc.unary_unary_rpc_method_handler(
servicer.Parse,
request_deserializer=parser__pb2.Request.FromString,
response_serializer=parser__pb2.Response.SerializeToString,
),
}
generic_handler = grpc.method_handlers_generic_handler(
'proto.Parser', rpc_method_handlers)
server.add_generic_rpc_handlers((generic_handler,))


# This class is part of an EXPERIMENTAL API.
class Parser(object):
"""Missing associated documentation comment in .proto file."""

@staticmethod
def Parse(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(request, target, '/proto.Parser/Parse',
parser__pb2.Request.SerializeToString,
parser__pb2.Response.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
Protobuf files for Go​
protoc -I proto/ proto/parser.proto --go_out=plugins=grpc:proto/
📎 parser.pb.go(8 KB)
Protobuf files for Java​
protoc --plugin=protoc-gen-grpc-java=./protoc-gen-grpc-java --grpc-java_out={out_dir} -I={input_dir} parser.proto
protoc --experimental_allow_proto3_optional -I {input_dir} --java_out={out_dir} parser.proto
ParserGrpc.java​
📎 ParserGrpc.java(10 KB)
ParserOuterClass.java​
📎 ParserOuterClass.java(34 KB)

Dependencies and versions:

com.google.protobuf:protobuf-java:3.19.4

io.grpc:grpc-all:1.51.0

Bundled package:

📎 te-logbus-grpc-customer-parse-1.0.0.jar(41.0 MB)
Was this page helpful?