跳到主要内容

插件配置 DEMO

最近更新 2026/10/03

logbus配置​

配置解析:

  • cmd 插件服务端启动命令
  • depend_sh 启动命令是否依赖bash环境,默认为false(python、java需要)
  • hand_shake 插件校验
  • hand_shake.protocol_version 与插件服务端定义须保持一致
  • hand_shake.magic_cookie_key 与插件服务端定义须保持一致
  • hand_shake.magic_cookie_value 与插件服务端定义须保持一致
{
"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"
}

插件开发​

grpc插件服务端接收数据格式​

数据批量处理,多条数据以双冒号分隔“::”,单条数据则拼接文件名称或者kafka的topic名称,以“|ta|”字符拼接。

example:

文件名称为:event.log

原始数据2条:{"user_id": "abc", "distinct_id": "123"} {"user_id": "def", "distinct_id": "456"}

grpc插件服务端接收到的数据为:event.log|ta|{"user_id": "abc", "distinct_id": "123"} ::event.log|ta|{"user_id": "def", "distinct_id": "456"}

grpc插件服务端返回数据格式​

以双冒号“::”拼接转换为标准ta格式的数据

example:

返回数据应为:

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

错误处理​

  1. 返回非正确json数据,则主程序会跳过错误行
  2. 返回非标准te数据,则数据上报会产生错误
  3. 返回第二个参数为error,则会进行无限重试,防止错误数据进入te集群。可搭配告警功能,提前预知错误信息。

程序demo【golang】​

protobuf定义参考附录,根据语言选择

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 自定义解析器
type Parser interface {
Parse([]byte) ([]byte, error)
}

// GRPCServer 插件服务端
type GRPCServer struct {
Impl Parser
}

type GRPCPlugin struct {
plugin.Plugin
Impl Parser
}

// GRPCClient GRPC客户端
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 服务端解析
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 服务端解析
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【python】​

校验

  • hand_shake.protocol_version 必填,值为1
  • hand_shake.magic_cookie_key 可忽略
  • hand_shake.magic_cookie_value 可忽略

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

注意这一行,按照格式输出到标准输出

1 固定,插件内定义

1 protocol_version

tcp 协议

address:port 地址

grpc 协议

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【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();
}
}

附录​

protobuf​

protobuf结构文件​
📎 parser.proto(1 KB)
python版本protobuf文件​
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)
go版本protobuf文件​
protoc -I proto/ proto/parser.proto --go_out=plugins=grpc:proto/
📎 parser.pb.go(8 KB)
java版本protobuf文件​
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)

依赖包及版本:

com.google.protobuf:protobuf-java:3.19.4

io.grpc:grpc-all:1.51.0

整合包:

📎 te-logbus-grpc-customer-parse-1.0.0.jar(41.0 MB)
这篇文档对你有帮助吗?