脚本之家,脚本语言编程技术及教程分享平台!
分类导航

Python|VBS|Ruby|Lua|perl|VBA|Golang|PowerShell|Erlang|autoit|Dos|bat|

服务器之家 - 脚本之家 - Golang - go实现grpc四种数据流模式

go实现grpc四种数据流模式

2022-09-14 14:06Jeff的技术栈 Golang

这篇文章主要为大家介绍了go实现grpc四种数据流模式,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步早日升职加薪

1. 什么是数据流

grpc中的stream,srteam顾名思义就是一种流,可以源源不断的推送数据,很适合传输一些大数据,或者服务端和客户端长时间数据交互,比如客户端可以向服务端订阅一个数据,服务端就可以利用stream,源源不断地推送数据。

底层还原成socket编程

2. grpc的四种数据流

1.简单模式

2.服务端数据流模式(Server-side streaming RPC)

3.客户端数据流模式(Client-side streaming RPC)

4.双向数据流模式(Bidirectional streaming RPC)

2.1 简单模式

  这种模式最为传统,即客户端发起一次请求,服务端响应一个数据,这和大家平时熟悉的RPC没有什么大的区别,上两篇中介绍此模式。

2.2 服务端数据流模式

  这种模式是客户端发起一次请求,服务端返回一段连续的数据流。典型的例子是客户端向服务端发送一个股票代码,服务端就把该股票的实时数据源源不断的返回给客户端

2.3 客户端数据流模式

  与服务端数据流模式相反,这次是客户端源源不断的向服务端发送数据流,而在发送结束后,由服务端返回一个响应。典型的例子是物联网终端向服务器报送数据。

2.4 双向数据流

  顾名思义,这是客户端和服务端都可以向对方发送数据流,这个时候双方的数据可以同时互相发送,也就是可以实现实时交互。典型的例子是聊天机器人。

3. 上代码

3.1 代码目录

go实现grpc四种数据流模式

3.2 编写stream.proto文件

stream是常量,写在哪一边,哪一边就是数据流

?
1
2
3
4
5
6
7
8
9
10
11
12
13
14
syntax = "proto3";
option go_package = "./;proto";
service Greeter {
    // 定义方法,stream是常量,流模式
    rpc ServerStream (StreamRequestData) returns (stream StreamResponseData);      //服务端流模式,拉消息
    rpc ClientStream (stream StreamRequestData) returns (StreamResponseData);      //客户端流模式,推消息
    rpc AllStream (stream StreamRequestData) returns (stream StreamResponseData);  //双向流模式,能推能拉
}
message StreamRequestData {
    string data = 1; //编号
}
message StreamResponseData {
    string data = 1; //编号
}

 生成go的protobuf文件命令:

cd到proto目录下

命令:protoc -I . hello.proto   --go_out=plugins=grpc:.

3.3 编写server文件

?
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
package main
import (
    "file_test/grpc_go_stream/proto"
    "fmt"
    "net"
    "sync"
    "time"
    "google.golang.org/grpc"
)
const port = 8082
type server struct{}
func (s *server) ServerStream(req *proto.StreamRequestData, res proto.Greeter_ServerStreamServer) error {
    i := 0
    for {
        i++
        //业务代码
        _ = res.Send(&proto.StreamResponseData{
            Data: fmt.Sprintf("这是发给%s的数据流", req.Data),
        })
        time.Sleep(time.Second * 1)
        if i > 10 {
            break
        }
    }
    return nil
}
func (s *server) ClientStream(cliStr proto.Greeter_ClientStreamServer) error {
    for {
        //业务代码
        res, err := cliStr.Recv()
        if err != nil {
            fmt.Println("本次客户端流数据发送完了:",err)
            break
        }
        fmt.Println("客户端发来消息:",res.Data)
    }
    return nil
}
func (s *server) AllStream(allStr proto.Greeter_AllStreamServer) error {
    wg:=sync.WaitGroup{}
    wg.Add(2)
    //接受客户端消息的协程
    go func() {
        defer wg.Done()
        for  {
            //业务代码
            res, err := allStr.Recv()
            if err != nil {
                fmt.Println("本次客户端流数据发送完了:",err)
                break
            }
            fmt.Println("收到客户端发来消息:",res.Data)
        }
    }()
    //发送消息给客户端的协程
    go func() {
        defer wg.Done()
        i := 0
        for {
            i++
            //业务代码
            _ = allStr.Send(&proto.StreamResponseData{
                Data: fmt.Sprintf("这是发给客户端的数据流"),
            })
            time.Sleep(time.Second * 1)
            if i > 10 {
                break
            }
        }
    }()
    wg.Wait()
    return nil
}
// 启动
func start() {
    // 1.实例化server
    g := grpc.NewServer()
    // 2.注册逻辑到server中
    proto.RegisterGreeterServer(g, &server{})
    // 3.启动server
    lis, err := net.Listen("tcp", "127.0.0.1:8082")
    if err != nil {
        panic("监听错误:" + err.Error())
    }
    err = g.Serve(lis)
    if err != nil {
        panic("启动错误:" + err.Error())
    }
}
func main() {
    start()
}

3.4 编写client文件

?
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
package main
import (
    "context"
    "file_test/grpc_go_stream/proto"
    "fmt"
    "sync"
    "time"
 
    "google.golang.org/grpc"
)
var rpc proto.GreeterClient
func serverStreamDemo()  {
    //服务端流模式
    res,err:=rpc.ServerStream(context.Background(),&proto.StreamRequestData{Data: "jeff"})
    if err != nil {
        panic("rpc请求错误:"+err.Error())
    }
    for  {
        data,err:=res.Recv() //
        if err != nil {
            fmt.Println("客户端发送完了:",err)
            return
        }
        fmt.Println("客户端返回数据流值:",data.Data)
    }
}
func clientStreamDemo()  {
    //客户端流模式
    cliStr, err := rpc.ClientStream(context.Background())
    if err != nil {
        panic("rpc请求错误:" + err.Error())
    }
    i := 0
    for {
        i++
        _ = cliStr.Send(&proto.StreamRequestData{
            Data: "jeff",
        })
        time.Sleep(time.Second * 1)
        if i > 10 {
            break
        }
    }
}
func clientAndServerStreamDemo()  {
    //双向流模式
    allStr, _ := rpc.AllStream(context.Background())
    wg := sync.WaitGroup{}
    wg.Add(1)
    //接受服务端消息的协程
    go func() {
        defer wg.Done()
        for {
            //业务代码
            res, err := allStr.Recv()
            if err != nil {
                fmt.Println("本次服务端流数据发送完了:", err)
                break
            }
            fmt.Println("收到服务端发来消息:", res.Data)
        }
    }()
    //发送消息给服务端的协程
    go func() {
        defer wg.Done()
        i := 0
        for {
            i++
            //业务代码
            _ = allStr.Send(&proto.StreamRequestData{
                Data: fmt.Sprintf("这是发给服务端的数据流"),
            })
            time.Sleep(time.Second * 1)
            if i > 10 {
                break
            }
        }
    }()
    wg.Wait()
}
// 启动
func start() {
    conn, err := grpc.Dial("127.0.0.1:8082", grpc.WithInsecure())
    if err != nil {
        panic("rpc连接错误:" + err.Error())
    }
    defer conn.Close()
    rpc = proto.NewGreeterClient(conn) //初始化
    serverStreamDemo() //服务端流模式
    clientStreamDemo()  //客户端流模式
    clientAndServerStreamDemo() // 双向流模式
}
func main() {
    start()
}

以上就是go实现grpc四种数据流模式的详细内容,更多关于go实现grpc流模式的资料请关注服务器之家其它相关文章!

原文链接:https://www.cnblogs.com/guyouyin123/p/16135335.html

延伸 · 阅读

精彩推荐
  • Golanggo嵌套匿名结构体的初始化详解

    go嵌套匿名结构体的初始化详解

    这篇文章主要介绍了go嵌套匿名结构体的初始化详解,具有很好的参考价值,希望对大家有所帮助。一起跟随小编过来看看吧...

    darkalliance4742021-02-28
  • GolangGo语言中的IO操作及Flag包的用法

    Go语言中的IO操作及Flag包的用法

    这篇文章介绍了Go语言中的IO操作及Flag包的用法,文中通过示例代码介绍的非常详细。对大家的学习或工作具有一定的参考借鉴价值,需要的朋友可以参考...

    奋斗的大橙子8952022-07-19
  • GolangGO中sync包自由控制并发示例详解

    GO中sync包自由控制并发示例详解

    这篇文章主要为大家介绍了GO中sync包自由控制并发示例详解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪...

    六号积极分子6532022-08-04
  • Golang在Go中使用JSON(附demo)

    在Go中使用JSON(附demo)

    Go开发人员经常需要处理JSON内容,本文主要介绍了在Go中使用JSON,文中通过示例代码介绍的非常详细,具有一定的参考价值,感兴趣的小伙伴们可以参考一...

    迪鲁宾7092022-09-05
  • Golanggolang通过context控制并发的应用场景实现

    golang通过context控制并发的应用场景实现

    这篇文章主要介绍了golang通过context控制并发的应用场景实现,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的...

    只是一个id3262020-06-02
  • Golang解决Go语言time包数字与时间相乘的问题

    解决Go语言time包数字与时间相乘的问题

    这篇文章主要介绍了Go语言time包数字与时间相乘的问题,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友可以参考下...

    西京刀客7752022-09-13
  • Golanggo语言实现一个最简单的http文件服务器实例

    go语言实现一个最简单的http文件服务器实例

    这篇文章主要介绍了go语言实现一个最简单的http文件服务器的方法,实例分析了Go语言操作http的技巧,具有一定参考借鉴价值,需要的朋友可以参考下 ...

    脚本之家4342020-04-22
  • Golanggolang 使用sort.slice包实现对象list排序

    golang 使用sort.slice包实现对象list排序

    这篇文章主要介绍了golang 使用sort.slice包实现对象list排序,对比sort跟slice两种排序的使用方式区别展开内容,需要的小伙伴可以参考一下...

    峰啊疯了10232022-09-09