Sign inSign up

hao474798383/msg-router

By hao474798383

•Updated over 1 year ago

kafka 到 web 端数据摆渡服务

Image
0

2.1K

hao474798383/msg-router repository overview

⁠msg-router kafka 消息-前端路由服务

⁠镜像名称:

hao474798383/msg-router:v2.3.3-jdk11

⁠1. 架构升级说明

message-router 2.x 版本升级说明:

  • 更新: java8 修改为使用 java11
  • 优化: 多端订阅, 多个 websocket 连接对相同管道的相同主题进行订阅的逻辑, 进行了大量优化工作
  • 优化: 对原有架构进行扩展性升级, 最大限度的兼容了 1.x 版本的交互协议
  • 新增: 鉴权服务, 在 websocket 握手时需要鉴权
  • 新增: 多管道订阅, 支持一个 websocket 连接订阅多个 kafka 管道的多个 topic
  • 新增: 多管道发布, 支持一个 websocket 连接对多个 kafka 管道的多个 topic 进行消息推送

⁠2. 使用说明

从 message-router 2.x 开始, web 端监听端口为 9102, websocket 端监听端口为 9002, 上下文均为 /

⁠2.1全局环境变量
参数说明是否必填默认值
KAFKA_HOSTSkafka 服务地址是kafka-0.kafka:9002
KAFKA_USERNAMEkafka 用户名否
KAFKA_PASSWORDkafka 密码否
KAFKA_IDENTIFY当前服务 ID 数字或英文组合 例如: qingzhi1, 多实例唯一多实例必须gomyck
KAFKA_POLL_SLEEP每次推送消息间隔 单位:毫秒否2000
KAFKA_POLL_RECORDS每次最多从 kafka 拉取多少条数据否1000
KAFKA_POLL_CHUNK_SIZE发送给前端时, 一次最多发多少条数据否1000

配置时, 应遵循下述公式(否则消费组会产生消息积压): KAFKA_POLL_RECORDS / KAFKA_POLL_SLEEP > 每秒进入 kafka 管道的数据条数

多实例指 多个 message-router 连接同一个 kafka, 如果不设置唯一 ID, 会导致共用一个 消费组(主题一致则经过算法得出的消费组一致) , 出现消费者闲置

⁠2.2 鉴权服务

鉴权服务默认关闭, 使用下述变量开启鉴权特性:

注意: 所有的环境变量都是非必填, 不填时, 取默认值!!!

变量名说明默认值
HYLINK_SECURE_ENABLED是否开启鉴权服务false
HYLINK_SECURE_ACCESS_KEY鉴权服务的 accessKeyhylink
HYLINK_SECURE_EXPIREDjwt token 的有效期,单位:小时12

鉴权服务开启时, 业务系统在使用管道前, 应由 业务系统服务端 请求 message-router 服务接口, 获取时限令牌, 并将令牌返回给前端页面使用, 接口如下:

{
  "result": true,
  "time": "2023-11-15 14:16:45",
  "message": "eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJtc2ctcm91dGVyIiwiaWF0IjoxNzAwMDI5MDA1LCJleHAiOjE3MDAwNzIyMDUsImNvcHlyaWdodCI6InRoaXMgdG9rZW4gaXMgZ2VuZXJhdGVkIGJ5IG1zZy1yb3V0ZXIsIElsbGVnYWwgdXNlIGlzIG5vdCBhbGxvd2VkISEhIn0.nFrKo0EX5fQA1_-a3CM-fczAJ1nOBINLql_mHwxX2X8",
  "event": "auth"
}
⁠2.3 websocket 握手

使用服务端获取的 token, 通过 websocket 握手进行鉴权

ws(s)://msg-router-ip:9002?token=eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJtc2ctcm91dGVyIiwiaWF0IjoxNzAwMDI5MDA1LCJleHAiOjE3MDAwNzIyMDUsImNvcHlyaWdodCI6InRoaXMgdG9rZW4gaXMgZ2VuZXJhdGVkIGJ5IG1zZy1yb3V0ZXIsIElsbGVnYWwgdXNlIGlzIG5vdCBhbGxvd2VkISEhIn0.nFrKo0EX5fQA1_-a3CM-fczAJ1nOBINLql_mHwxX2X8

⁠2.4 管道创建指令

如果使用message-router 默认的 kafka 管道, 不需要调用此命令!!!!!

message-router 在启动时, 配置了默认的 kafka 管道, 当需要使用其他 kafka 管道时, 可通过调用此指令, 创建新的 kafka 管道

该指令返回管道 ID, 用于后续指令调用中使用

入参:

{
  "event": "init_broker",
  "bootstrapServers": "192.168.3.71:9094",
  "username": "kafka 没有鉴权不传, 否则按实际传参",
  "password": "kafka 没有鉴权不传, 否则按实际传参"
}

返回:

{
    "event": "init_broker",
    "result": true,
    "time": "2023-11-15 16:11:16",
    "message": "05D08DEC8E86E339076213357BEE6D2A"
}
⁠2.5 消息订阅指令(兼容 1.x)

使用下面的参数进行订阅

注意: 因为服务端内置了默认管道, 所以只有 event 和 topics 两个参数是必填的, 其他都为非必填参数

{
    "event": "subscribe",  
    "brokerId": "05D08DEC8E86E339076213357BEE6D2A",
    "consumeType": "batch",
    "debug": "false",
    "beginTime": "2023-01-01 11:11:11",
    "saveGroup": "false",
    "topics": ["msg-router", "xxx.*"] 
}
参数说明是否必填默认值
event调用事件, 固定为 subscribe是subscribe
brokerId2.4 管道创建指令返回的 brokerId, 不传则使用默认的管道否subscribe
consumeTypebatch 批量接收, 返回的是数组, 前端统一处理, 节省前端算力
single 单条接收, 返回的是对象, 与 1.x 版本消息返回形式一致
否batch
debugfalse 仅返回 kafka 管道 ID 与 kafka 主题中的 message
true 响应中返回更多调试信息
否false
saveGroupfalse 不保存消费组, 每次都从当前时间开始消费
true 保存消费组, 下次重连可从上次断开位置消费
否false
beginTime从指定时间开始消费, 格式 yyyy-MM-dd HH:mm:ss否
topics字符串数组,支持正则表达式 例如: ["msg-router", "xxxx"]是

batch 返回:

[{
    "brokerId": "911BEB2D0CEC8CFDEC710C51AB32ECFA",
    "value": {
        "XXX": "ASDASD",
        "DDDD": "node-red 消息, 发至 默认 Wed Nov 15 2023 16:17:42 GMT+0800 (China Standard Time)"
    }
}, {
    "brokerId": "911BEB2D0CEC8CFDEC710C51AB32ECFA",
    "value": {
        "XXX": "ASDASD",
        "DDDD": "node-red 消息, 发至 默认 Wed Nov 15 2023 16:17:43 GMT+0800 (China Standard Time)"
    }
}]

single 返回:

{
    "brokerId": "911BEB2D0CEC8CFDEC710C51AB32ECFA",
    "value": {
        "XXX": "ASDASD",
        "DDDD": "node-red 消息, 发至 默认 Wed Nov 15 2023 16:16:36 GMT+0800 (China Standard Time)"
    }
}
⁠2.6 消息发布指令
{
    "event": "publish",
    "brokerId": "05D08DEC8E86E339076213357BEE6D2A",
    "topics": ["msg-router", "msg-router2"],
    "message": "测试消息"
}
参数说明是否必填默认值
event调用事件, 固定为 publish是subscribe
brokerId2.4 管道创建指令返回的 brokerId, 不传则使用默认的管道否subscribe
message推送到 kafka 的消息是
topics字符串数组,必须是精确匹配是

返回:

{
    "result": true,
    "time": "2023-11-15 16:22:45",
    "message": "success",
    "event": "publish"
}

⁠更新说明

v2.1.2-jdk11: 新增manticoreSearch 客户端订阅组件, 订阅指令新增两项参数

POST http://localhost:9102/subscribe⁠

{
  "event": "subscribe",
  "saveGroup": "false",
  "topics": ["test"],
  "url": "http://192.168.104.60:9308/insert",
  "target": "Manticoresearch"
}

v2.2.1-jdk11: 新增订阅参数: filterInfo, 响应实体的 json 中, 对应 key 如果包含指定的值, 则返回

{
    "filterInfo": {
        "userId": "xxx",
        "other": "zzz"
    }  
}

v2.2.2-jdk11:
1.新增订阅参数: jsonPathFilterInfo, 响应实体的 json 中, 如果符合 jsonpath的规则, 则返回
2.提升了 springboot 的版本, 修复了 BUG

{
    "jsonPathFilterInfo": {
        "$[1:]..xxx": "123"
    }  
}

v2.3.0-jdk11 (重要):
1.修复了在订阅相同主题时, 批量消费和单条消费不能同时存在的 bug

v2.3.3-jdk11 (重要):

  1. 修复了客户端之间会话污染

Tag summary

Content type

Image

Digest

sha256:2ad9ef454…

Size

133.6 MB

Last updated

over 1 year ago

docker pull hao474798383/msg-router