边缘计算不仅仅是将应用部署在边缘,并对其进行自动化的监控和运维。在许多应用场景里,边缘和云上应用需要进行特定的消息传输、数据交换等,以完成边云协同的业务处理。例如,用户需要从云端发送命令至边缘的应用来触发特定的业务,或者边缘设备需要将采集的业务信息上传至云端处理。
KubeEdge v1.6 版本增加了自定义边云消息传输的支持,用户可以根据场景,借助 Rule 和 RuleEndpoint 两个新增API来自定义的边云消息传输设置,为需要边云通信的业务组件或第三方插件屏蔽底层网络环境差异。
Router Manager
KubeEdge 借助 Kubernetes CRD 和路由器模块支持路由管理。用户可以通过mqtt代理在云和边缘之间传递其自定义消息。
使用场景: 用于用户控制数据的传递; 不适合大数据传送; 一次传送的数据大小限制为12MB。
kubeedge 升级到 1.6 版本貌似默认没有开启 router ,需要手动创建 router crd 和开启 router 功能。

修改 cloudcore 配置,开启 router 功能。如果是新环境,直接开启即可。
1# /etc/kubeedge/config/cloudcore.yaml 2 router: 3 enable: true 4 address: 0.0.0.0 5 port: 9443 6 restTimeout: 60 7 syncController: 8 enable: true 9 10# kubectl get crds | grep kubeedge 11clusterobjectsyncs.reliablesyncs.kubeedge.io 2021-03-08T07:53:39Z 12devicemodels.devices.kubeedge.io 2021-03-08T07:53:39Z 13devices.devices.kubeedge.io 2021-03-08T07:53:39Z 14objectsyncs.reliablesyncs.kubeedge.io 2021-03-08T07:53:39Z 15ruleendpoints.rules.kubeedge.io 2021-03-08T08:04:59Z 16rules.rules.kubeedge.io 2021-03-08T08:04:59Z
云端到边缘端
通过调用 cloudcore api 将控制边缘端应用的消息从云传递到边缘,边缘端的应用接收消息后,开始或停止收集树莓派的系统日志,收集到的日志发布到 mq ,日志可以通过 kuiper 的规则处理,把需要的日志上传到云端 emqx 。
环境
1# kubectl get nodes -o wide 2NAME STATUS ROLES AGE VERSION INTERNAL-IP EXTERNAL-IP OS-IMAGE KERNEL-VERSION CONTAINER-RUNTIME 3k8s-test-master01 Ready master 21d v1.20.2 172.31.250.220 <none> CentOS Linux 8 (Core) 4.18.0-193.28.1.el8_2.x86_64 cri-o://1.20.0 4k8s-test-node01 Ready node 21d v1.20.2 172.31.250.221 <none> Ubuntu 20.04.1 LTS 5.4.0-58-generic containerd://1.3.3-0ubuntu2.2 5k8s-test-node02 Ready node 21d v1.20.2 172.31.250.222 <none> openSUSE Leap 15.2 5.3.18-lp152.57-default cri-o://1.17.3 6kubeedge-raspberrypi01 Ready agent,edge 2d2h v1.19.3-kubeedge-v1.6.0 192.168.13.102 <none> Raspbian GNU/Linux 10 (buster) 5.4.51-v7l+ docker://19.3.13
1.创建云端到边缘端的路由
1# cat create-ruleEndpoint-rest.yaml 2apiVersion: rules.kubeedge.io/v1 3kind: RuleEndpoint 4metadata: 5 name: my-rest 6 labels: 7 description: test 8spec: 9 ruleEndpointType: "rest" 10 properties: {} 11--- 12# cat create-ruleEndpoint-eventbus.yaml 13apiVersion: rules.kubeedge.io/v1 14kind: RuleEndpoint 15metadata: 16 name: my-eventbus 17 labels: 18 description: test 19spec: 20 ruleEndpointType: "eventbus" 21 properties: {} 22--- 23# cat create-rule-rest-eventbus.yaml 24apiVersion: rules.kubeedge.io/v1 25kind: Rule 26metadata: 27 name: my-rule 28 labels: 29 description: test 30spec: 31 source: "my-rest" 32 sourceResource: {"path":"/a"} 33 target: "my-eventbus" 34 targetResource: {"topic":"test"}
2.写一个简单的程序,用于收集 Linux 日志
1// messagePubHandler 订阅 mq 消息,启动或停止日志采集 2var messagePubHandler mqtt.MessageHandler = func(client mqtt.Client, msg mqtt.Message) { 3 4 log.Printf("Received message: %v messageid: %d from topic: %s\n", msg, msg.MessageID(), msg.Topic()) 5 controlMessage := string(msg.Payload()) 6 log.Println(controlMessage) 7 if controlMessage == "start" { 8 fileName := "/var/log/syslog" 9 config := tail.Config{ 10 ReOpen: true, // 重新打开 11 Follow: true, // 是否跟随 12 Location: &tail.SeekInfo{Offset: S.offset, Whence: 0}, // 从文件的哪个地方开始读 13 MustExist: true, // 文件不存在报错 14 Poll: true, 15 } 16 17 tails, err := tail.TailFile(fileName, config) 18 if err != nil { 19 log.Println("tail file failed, err:", err) 20 return 21 } 22 go logsStream(tails, S.mqttClient) 23 24 } else if controlMessage == "stop" { 25 S.StopCh <- string(msg.Payload()) 26 } else { 27 log.Printf("Unknown message : %s", controlMessage) 28 } 29 30}

3.在 Kuiper 创建日志流
1/kuiper # bin/kuiper create stream logs '(month string, days string,times string ,hostname string,kinds string,logs string,) WITH (FORMAT="JSON", DATASOURCE="demo")'; 2Connecting to 127.0.0.1:20498... 3Stream logs is created. 4/kuiper # bin/kuiper show streams 5Connecting to 127.0.0.1:20498... 6logs
4.创建 Kuiper 规则,把日志转发到云端 emqx
1cat > /tmp/rule.yaml <<EOF 2{ 3 "id": "rule1", 4 "sql": "select * from logs ", 5 "actions": [ 6 { 7 "log": {} 8 }, 9 { 10 "mqtt": { 11 "server": "tcp://${云端emqx IP}:1883", 12 "topic": "demoSink" 13 } 14 } 15 ] 16} 17EOF 18/kuiper # bin/kuiper create rule rule1 -f /tmp/rule.yaml
5.调用云端的 cloudcore rest api 将消息发送到边缘端
1# URL: http://{rest_endpoint}/{node_name}/{namespace}/{path} 2[root@k8s-test-master01 ~]# curl -X POST --data 'start' http://127.0.0.1:9443/kubeedge-raspberrypi01/default/a 3message delivered
6.边缘端的程序接收到消息,开始收集日志

7.订阅云端的 emqx 消息,验证是否有日志发送过来

边缘端到云端
1.创建边缘端到云端的路由
1# create-ruleEndpoint-rest.yaml 2apiVersion: rules.kubeedge.io/v1 3kind: RuleEndpoint 4metadata: 5 name: my-rest 6 labels: 7 description: test 8spec: 9 ruleEndpointType: "rest" 10 properties: {} 11--- 12# create-ruleEndpoint-eventbus.yaml 13apiVersion: rules.kubeedge.io/v1 14kind: RuleEndpoint 15metadata: 16 name: my-eventbus 17 labels: 18 description: test 19spec: 20 ruleEndpointType: "eventbus" 21 properties: {} 22--- 23#create-rule-eventbus-rest.yaml 24apiVersion: rules.kubeedge.io/v1 25kind: Rule 26metadata: 27 name: my-rule-eventbus-rest 28 labels: 29 description: test 30spec: 31 source: "my-eventbus" 32 sourceResource: {"topic": "test","node_name": "kubeedge-raspberrypi01"} 33 target: "my-rest" 34 targetResource: {"resource":"http://172.31.250.220:8088/api/v1/msg"}
resource 是云端应用的 api 地址
2.写一个简单的 API 接口
1package main 2 3import( 4 "github.com/kataras/iris/v12" 5 "github.com/kataras/iris/v12/middleware/logger" 6 "github.com/kataras/iris/v12/middleware/recover" 7) 8 9type Msg struct{ 10 EdgeMsg string `json:"edgemsg"` 11} 12 13func main(){ 14 app := iris.New() 15 c := &Msg{} 16 app.Logger().SetLevel("debug") 17 18 app.Use(recover.New()) 19 app.Use(logger.New()) 20 21 resAPI := app.Party("/api/v1") 22 resAPI.Post("/msg", func(ctx iris.Context){ 23 if err := ctx.ReadJSON(c); err != nil{ 24 panic(err.Error()) 25 }else{ 26 ctx.JSON(c) 27 } 28 }) 29 30 resAPI.Get("/msg", func(ctx iris.Context){ 31 ctx.Writef("Received: %v\n", c.EdgeMsg) 32 }) 33 34 app.Run(iris.Addr(":8088"), iris.WithoutServerError(iris.ErrServerClosed)) 35}

目前接口没有数据

3.在边缘端的将自定义的消息发布到边缘节点的 MQTT 代理

1mosquitto_pub -t 'default/test' -d -m '{"edgemsg":"msgtocloud"}'
4.再访问云端 API 接口,这时已经获取到了从边缘端发送的消息。

结束
因为目前边云自定义消息传输的不适合大数据传送,一次传送的数据大小限制为12MB 和单向 REST 的局限性,目前使用场景还是相对简单,可能更多是用于用户控制数据的传递,比如控制边缘终端设备的启停、边缘端向云端汇总边缘终端设备的在线或离线状态等。
社区也计划在下一个版本中优化和扩展该功能特性。 参考文档 https://kubeedge.io/en/docs/developer/custom_message_deliver/\
感兴趣的读者可以关注下微信号

