【问题标题】:gRPC keeping response streams open for subscriptionsgRPC 保持响应流对订阅开放
【发布时间】:2020-10-07 17:48:03
【问题描述】:

我尝试定义一个 gRPC 服务,客户端可以订阅接收广播消息,也可以发送它们。

syntax = "proto3";

package Messenger;

service MessengerService {
    rpc SubscribeForMessages(User) returns (stream Message) {}
    rpc SendMessage(Message) returns (Close) {}
}

message User {
    string displayName = 1;
}

message Message {
    User from = 1;
    string message = 2;
}

message Close {}

我的想法是,当客户端请求订阅消息时,响应流将被添加到响应流集合中,当发送消息时,消息通过所有响应流发送。

但是,当我的服务器尝试写入响应流时,我收到异常 System.InvalidOperationException: 'Response stream has already been completed.'

有没有办法告诉服务器保持流打开以便可以通过它们发送新消息?或者这不是 gRPC 的设计目的,应该使用不同的技术?

最终目标服务将允许通过用不同语言(C#、Java 等)编写的不同客户端进行多种类型的订阅(可能是新消息、天气更新等)。不同的语言部分主要是我选择 gRPC 来尝试这个的原因,虽然我打算用 C# 编写服务器。


实现示例

using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Grpc.Core;
using Messenger;

namespace SimpleGrpcTestStream
{
    /*
    Dependencies
Install-Package Google.Protobuf
Install-Package Grpc
Install-Package Grpc.Tools
Install-Package System.Interactive.Async
Install-Package System.Linq.Async

    */
    internal static class Program
    {
        private static void Main()
        {
            var messengerServer = new MessengerServer();
            messengerServer.Start();

            var channel = Common.GetNewInsecureChannel();
            var client = new MessengerService.MessengerServiceClient(channel);
            var clientUser = Common.GetUser("Client");
            var otherUser = Common.GetUser("Other");

            var cancelClientSubscription = AddCancellableMessageSubscription(client, clientUser);
            var cancelOtherSubscription = AddCancellableMessageSubscription(client, otherUser);

            client.SendMessage(new Message { From = clientUser, Message_ = "Hello" });
            client.SendMessage(new Message { From = otherUser, Message_ = "World" });
            client.SendMessage(new Message { From = clientUser, Message_ = "Whoop" });

            cancelClientSubscription.Cancel();
            cancelOtherSubscription.Cancel();
            channel.ShutdownAsync().Wait();
            messengerServer.ShutDown().Wait();
        }

        private static CancellationTokenSource AddCancellableMessageSubscription(
            MessengerService.MessengerServiceClient client,
            User user)
        {
            var cancelMessageSubscription = new CancellationTokenSource();

            var messages = client.SubscribeForMessages(user);

            var messageSubscription = messages
                .ResponseStream
                .ToAsyncEnumerable()
                .Finally(() => messages.Dispose());

            messageSubscription.ForEachAsync(
                message => Console.WriteLine($"New Message: {message.Message_}"),
                cancelMessageSubscription.Token);

            return cancelMessageSubscription;
        }
    }

    public static class Common
    {
        private const int Port = 50051;

        private const string Host = "localhost";

        private static readonly string ChannelAddress = $"{Host}:{Port}";

        public static User GetUser(string name) => new User { DisplayName = name };

        public static readonly User ServerUser = GetUser("Server");

        public static readonly Close EmptyClose = new Close();

        public static Channel GetNewInsecureChannel() => new Channel(ChannelAddress, ChannelCredentials.Insecure);

        public static ServerPort GetNewInsecureServerPort() => new ServerPort(Host, Port, ServerCredentials.Insecure);
    }

    public sealed class MessengerServer : MessengerService.MessengerServiceBase
    {
        private readonly Server _server;

        public MessengerServer()
        {
            _server = new Server
            {
                Ports = { Common.GetNewInsecureServerPort() },
                Services = { MessengerService.BindService(this) },
            };
        }

        public void Start()
        {
            _server.Start();
        }

        public async Task ShutDown()
        {
            await _server.ShutdownAsync().ConfigureAwait(false);
        }

        private readonly ConcurrentDictionary<User, IServerStreamWriter<Message>> _messageSubscriptions = new ConcurrentDictionary<User, IServerStreamWriter<Message>>();

        public override async Task<Close> SendMessage(Message request, ServerCallContext context)
        {
            await Task.Run(() =>
            {
                foreach (var (_, messageStream) in _messageSubscriptions)
                {
                    messageStream.WriteAsync(request);
                }
            }).ConfigureAwait(false);

            return await Task.FromResult(Common.EmptyClose).ConfigureAwait(false);
        }

        public override async Task SubscribeForMessages(User request, IServerStreamWriter<Message> responseStream, ServerCallContext context)
        {
            await Task.Run(() =>
            {
                responseStream.WriteAsync(new Message
                {
                    From = Common.ServerUser,
                    Message_ = $"{request.DisplayName} is listening for messages!",
                });
                _messageSubscriptions.TryAdd(request, responseStream);
            }).ConfigureAwait(false);
        }
    }

    public static class AsyncStreamReaderExtensions
    {
        public static IAsyncEnumerable<T> ToAsyncEnumerable<T>(this IAsyncStreamReader<T> asyncStreamReader)
        {
            if (asyncStreamReader is null) { throw new ArgumentNullException(nameof(asyncStreamReader)); }

            return new ToAsyncEnumerableEnumerable<T>(asyncStreamReader);
        }

        private sealed class ToAsyncEnumerableEnumerable<T> : IAsyncEnumerable<T>
        {
            public IAsyncEnumerator<T> GetAsyncEnumerator(CancellationToken cancellationToken = default)
                => new ToAsyncEnumerator<T>(_asyncStreamReader, cancellationToken);

            private readonly IAsyncStreamReader<T> _asyncStreamReader;

            public ToAsyncEnumerableEnumerable(IAsyncStreamReader<T> asyncStreamReader)
            {
                _asyncStreamReader = asyncStreamReader;
            }

            private sealed class ToAsyncEnumerator<TEnumerator> : IAsyncEnumerator<TEnumerator>
            {
                public TEnumerator Current => _asyncStreamReader.Current;

                public async ValueTask<bool> MoveNextAsync() => await _asyncStreamReader.MoveNext(_cancellationToken);

                public ValueTask DisposeAsync() => default;

                private readonly IAsyncStreamReader<TEnumerator> _asyncStreamReader;
                private readonly CancellationToken _cancellationToken;

                public ToAsyncEnumerator(IAsyncStreamReader<TEnumerator> asyncStreamReader, CancellationToken cancellationToken)
                {
                    _asyncStreamReader = asyncStreamReader;
                    _cancellationToken = cancellationToken;
                }
            }
        }
    }
}

【问题讨论】:

    标签: c# .net-core grpc


    【解决方案1】:

    您遇到的问题是由于MessengerServer.SubscribeForMessages 立即返回。一旦该方法返回,流就关闭了。

    你需要一个类似的实现来保持流的活跃:

    public class MessengerService : MessengerServiceBase
    {
        private static readonly ConcurrentDictionary<User, IServerStreamWriter<Message>> MessageSubscriptions =
            new Dictionary<User, IServerStreamWriter<Message>>();
    
        public override async Task SubscribeForMessages(User request, IServerStreamWriter<ReferralAssignment> responseStream, ServerCallContext context)
        {
            if (!MessageSubscriptions.TryAdd(request))
            {
                // User is already subscribed
                return;
            }
    
            // Keep the stream open so we can continue writing new Messages as they are pushed
            while (!context.CancellationToken.IsCancellationRequested)
            {
                // Avoid pegging CPU
                await Task.Delay(100);
            }
    
            // Cancellation was requested, remove the stream from stream map
            MessageSubscriptions.TryRemove(request);
        }
    }
    

    就退订/取消而言,有两种可能的方法:

    1. 客户端可以保持CancellationToken,并在想要断开连接时调用Cancel()
    2. 服务器可以保留CancellationToken,然后您可以通过Tuple 或类似方法将其与IServerStreamWriter 一起存储在MessageSubscriptions 字典中。然后,您可以在服务器上引入一个Unsubscribe 方法,该方法通过User 查找CancellationToken 并在其服务器端调用Cancel

    【讨论】:

    • 感谢您的回答。发现它真的很有用,并让我知道了我想如何实现它。在赏金结束时,如果没有更好的方法出现,我会将赏金奖励给你
    【解决方案2】:

    类似于Jon Halliday's 答案,可以使用无限长的Task.Delay(-1) 并传递上下文的取消令牌。

    当任务被取消时,可以使用 try catch 移除结束服务器的响应流。

    public override async Task SubscribeForMessages(User request, IServerStreamWriter<Message> responseStream, ServerCallContext context)
    {
        if (_messageSubscriptions.ContainsKey(request))
        {
            return;
        }
    
        await responseStream.WriteAsync(new Message
        {
            From = Common.ServerUser,
            Message_ = $"{request.DisplayName} is listening for messages!",
        }).ConfigureAwait(false);
    
        _messageSubscriptions.TryAdd(request, responseStream);
    
        try
        {
            await Task.Delay(-1, context.CancellationToken);
        }
        catch (TaskCanceledException)
        {
            _messageSubscriptions.TryRemove(request, out _);
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-12-05
      • 2019-09-26
      • 2018-06-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多