gRPC - 服务器流式 RPC
现在让我们讨论一下使用 gRPC 通信时服务器流式传输的工作原理。在这种情况下,客户端将搜索具有给定作者的书籍。假设服务器需要一些时间来浏览所有书籍。服务器不会等待浏览完所有书籍后再提供所有书籍,而是以流式传输的方式提供书籍,即,只要找到一本就提供。
.proto 文件
首先让我们在 common_proto_files 中定义 bookstore.proto 文件 −
syntax = "proto3"; option java_package = "com.tp.bookstore"; service BookStore { rpc first (BookSearch) returns (stream Book) {} } message BookSearch { string name = 1; string author = 2; string genre = 3; } message Book { string name = 1; string author = 2; int32 price = 3; }
以下块表示服务的名称"BookStore"和可调用的函数名称"searchByAuthor"。"searchByAuthor"函数接受类型为"BookSearch"的输入并返回类型为"Book"的流。因此,实际上,我们让客户端搜索标题并返回与查询的作者匹配的一本书。
service BookStore { rpc searchByAuthor (BookSearch) returns (stream Book) {} }
Now let us look at these types.
message BookSearch { string name = 1; string author = 2; string genre = 3; }
这里,我们定义了 BookSearch,它包含几个属性,如 name、author 和 genre。客户端应该将 "BookSearch" 类型的对象发送到服务器。
message Book { string name = 1; string author = 2; int32 price = 3; }
我们还定义了,给定一个 "BookSearch",服务器将返回一个 "Book" 流,其中包含书籍属性以及书籍价格。服务器应该发送一个"Book"流。
请注意,我们已经完成了 Maven 设置,可以自动生成我们的类文件以及我们的 RPC 代码。所以,现在我们可以简单地编译我们的项目 −
mvn clean install
这应该会自动生成我们使用 gRPC 所需的源代码。源代码将放在 −
Protobuf class code: target/generated-sources/protobuf/java/com.tp.bookstore Protobuf gRPC code: target/generated-sources/protobuf/grpc-java/com.tp.bookstore
设置 gRPC 服务器
现在我们已经定义了包含函数定义的 proto 文件,让我们设置一个可以调用这些函数的服务器。
让我们编写服务器代码来提供上述功能并将其保存在 com.tp.bookstore.BookeStoreServerStreaming.java −
示例
package com.tp.bookstore; import io.grpc.Server; import io.grpc.ServerBuilder; import io.grpc.stub.StreamObserver; import java.io.IOException; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.concurrent.TimeUnit; import java.util.logging.Logger; import java.util.stream.Collectors; import com.tp.bookstore.BookStoreOuterClass.Book; import com.tp.bookstore.BookStoreOuterClass.BookSearch; public class BookeStoreServerUnary { private static final Logger logger = Logger.getLogger(BookeStoreServerrStreaming.class.getName()); static Map<String, Book> bookMap = new HashMap<>(); static { bookMap.put("Great Gatsby", Book.newBuilder().setName("Great Gatsby") .setAuthor("Scott Fitzgerald") .setPrice(300).build()); bookMap.put("To Kill MockingBird", Book.newBuilder().setName("To Kill MockingBird") .setAuthor("Harper Lee") .setPrice(400).build()); bookMap.put("Passage to India", Book.newBuilder().setName("Passage to India") .setAuthor("E.M.Forster") .setPrice(500).build()); bookMap.put("The Side of Paradise", Book.newBuilder().setName("The Side of Paradise") .setAuthor("Scott Fitzgerald") .setPrice(600).build()); bookMap.put("Go Set a Watchman", Book.newBuilder().setName("Go Set a Watchman") .setAuthor("Harper Lee") .setPrice(700).build()); } private Server server; private void start() throws IOException { int port = 50051; server = ServerBuilder.forPort(port) .addService(new BookStoreImpl()).build().start(); logger.info("Server started, listening on " + port); Runtime.getRuntime().addShutdownHook(new Thread() { @Override public void run() { System.err.println("Shutting down gRPC server"); try { server.shutdown().awaitTermination(30, TimeUnit.SECONDS); } catch (InterruptedException e) { e.printStackTrace(System.err); } } }); } public static void main(String[] args) throws IOException, InterruptedException { final BookeStoreServerUnary greetServer = new BookeStoreServerUnary(); greetServer.start(); greetServer.server.awaitTermination(); } static class BookStoreImpl extends BookStoreGrpc.BookStoreImplBase { @Override public void searchByAuthor(BookSearch searchQuery, StreamObserver<Book> responseObserver) { logger.info("Searching for book with author: " + searchQuery.getAuthor()); for (Entry<String, Book> bookEntry : bookMap.entrySet()) { try { logger.info("Going through more books...."); Thread.sleep(5000); } catch (InterruptedException e) { e.printStackTrace(); } if(bookEntry.getValue().getAuthor().startsWith(searchQuery.getAuthor())){ logger.info("Found book with required author: " + bookEntry.getValue().getName()+ ". Sending...."); responseObserver.onNext(bookEntry.getValue()); } } responseObserver.onCompleted(); } } }
上述代码在指定端口启动一个 gRPC 服务器,并提供我们在 proto 文件中编写的功能和服务。让我们看一下上面的代码 −
从 main 方法开始,我们在指定端口创建一个 gRPC 服务器。
但在启动服务器之前,我们为服务器分配了我们想要运行的服务,即在我们的例子中为 BookStore 服务。
为此,我们需要将服务实例传递给服务器,因此我们继续创建一个服务实例,即在我们的例子中为 BookStoreImpl
服务实例需要提供 .proto 文件 中存在的方法/函数的实现,即在我们的例子中为 searchByAuthor 方法。
该方法需要一个在 .proto 文件中定义的类型的对象,即,对于我们来说,BookSearch
请注意,我们添加了一个 sleep 来模拟搜索所有书籍的操作。在流式传输的情况下,服务器不会等待所有搜索到的书籍都可用。它通过使用 onNext() 调用在书籍可用时立即返回书籍。
当服务器完成请求后,它会通过调用 onCompleted() 关闭通道。
最后,我们还有一个关闭钩子,以确保在执行完代码后干净地关闭服务器。
设置 gRPC 客户端
现在我们已经编写了服务器的代码,让我们设置一个可以调用这些函数的客户端。
让我们编写客户端代码来调用上述函数并将其保存在 com.tp.bookstore.BookStoreClientServerStreamingBlocking.java −
示例
package com.tp.bookstore; import io.grpc.Channel; import io.grpc.ManagedChannel; import io.grpc.ManagedChannelBuilder; import io.grpc.StatusRuntimeException; import java.util.Iterator; import java.util.concurrent.TimeUnit; import java.util.logging.Level; import java.util.logging.Logger; import com.tp.bookstore.BookStoreOuterClass.Book; import com.tp.bookstore.BookStoreOuterClass.BookSearch; import com.tp.greeting.GreeterGrpc; import com.tp.greeting.Greeting.ServerOutput; import com.tp.greeting.Greeting.ClientInput; public class BookStoreClientServerStreamingBlocking { private static final Logger logger = Logger.getLogger(BookStoreClientServerStreamingBlocking.class.getName()); private final BookStoreGrpc.BookStoreBlockingStub blockingStub; public BookStoreClientServerStreamingBlocking(Channel channel) { blockingStub = BookStoreGrpc.newBlockingStub(channel); } public void getBook((String author) { logger.info("Querying for book with author: " + author); BookSearch request = BookSearch.newBuilder()..setAuthor(author).build(); Iterator<Book> response; try { response = blockingStub.searchByAuthor(request); while(response.hasNext()) { logger.info("Found book: " + response.next()); } } catch (StatusRuntimeException e) { logger.log(Level.WARNING, "RPC failed: {0}", e.getStatus()); return; } } public static void main(String[] args) throws Exception { String authorName = args[0]; String serverAddress = "localhost:50051"; ManagedChannel channel = ManagedChannelBuilder.forTarget(serverAddress) .usePlaintext() .build(); try { BookStoreClientServerStreamingBlocking client = new BookStoreClientUnaryBlocking(channel); client.getBook(authorName); } finally { channel.shutdownNow().awaitTermination(5, TimeUnit.SECONDS); } } }
上述代码在指定端口启动 gRPC 服务器,并提供我们在 proto 文件中编写的功能和服务。让我们来看看上面的代码 −
从 main 方法开始,我们接受一个参数,即我们要搜索的书的 标题。
我们设置了一个与服务器进行 gRPC 通信的通道。
然后,我们使用该通道创建一个 阻塞存根。在这里我们选择了我们计划调用其功能的服务"BookStore"。
然后,我们只需创建 .proto 文件中定义的预期输入,即我们的示例中的 BookSearch,并添加我们希望服务器搜索的标题。
我们最终进行调用并获取有效图书的迭代器。当我们迭代时,我们会获得服务器提供的相应书籍。
最后,我们关闭通道以避免任何资源泄漏。
所以,这就是我们的客户端代码。
客户端服务器调用
总而言之,我们想要做的是以下 −
启动 gRPC 服务器。
客户端向服务器查询具有给定作者的书籍。
服务器在其商店中搜索书籍,这是一个耗时的过程。
只要找到符合给定条件的书籍,服务器就会做出响应。服务器不会等待所有有效书籍都可用。它会在找到一本后立即发送输出。然后重复该过程。
现在,我们已经定义了 proto 文件,编写了服务器和客户端代码,让我们继续执行此代码并查看实际操作。
要运行代码,请启动两个 shell。通过执行以下命令在第一个 shell 上启动服务器 −
java -cp . arget\grpc-point-1.0.jar com.tp.bookstore.BookeStoreServerStreaming
我们将看到以下输出 −
输出
Jul 03, 2021 10:37:21 PM com.tp.bookstore.BookeStoreServerStreaming start INFO: Server started, listening on 50051
以上输出表示服务器已经启动。
现在,让我们启动客户端。
java -cp . arget\grpc-point-1.0.jar com.tp.bookstore.BookStoreClientServerStreamingBlocking "Har"
我们将看到以下输出 −
输出
Jul 03, 2021 10:40:31 PM com.tp.bookstore.BookStoreClientServerStreamingBlocking getBook INFO: Querying for book with author: Har Jul 03, 2021 10:40:37 PM com.tp.bookstore.BookStoreClientServerStreamingBlocking getBook INFO: Found book: name: "Go Set a Watchman" author: "Harper Lee" price: 700 Jul 03, 2021 10:40:42 PM com.tp.bookstore.BookStoreClientServerStreamingBlocking getBook INFO: Found book: name: "To Kill MockingBird" author: "Harper Lee" price: 400
因此,如我们所见,客户端能够通过查询书名来获取书籍详细信息。但更重要的是,客户端在不同的时间戳获取了第一本书和第二本书,也就是说,间隔近 5 秒。