C#定时拉取GCP Pub/Sub消息时出现Connection reset by peer异常求助
Google Cloud Pub/Sub 连接重置异常问题排查与解决
问题描述
我在C#应用中使用Google Cloud Pub/Sub,通过定时任务每5秒从主题拉取消息。运行一段时间后,遇到如下异常:
System.IO.IOException: Unable to read data from the transport connection: Connection reset by peer. ---> System.Net.Sockets.SocketException (104): Connection reset by peer --- End of inner exception stack trace --- at Google.Cloud.PubSub.V1.SubscriberClientImpl.SingleChannel.HandleRpcFailure(Exception e) at Google.Cloud.PubSub.V1.SubscriberClientImpl.SingleChannel.HandlePullMoveNext(Task initTask) at Google.Cloud.PubSub.V1.SubscriberClientImpl.SingleChannel.StartAsync() at Google.Cloud.PubSub.V1.Tasks.ForwardingAwaiter.GetResult() at Google.Cloud.PubSub.V1.Tasks.Extensions.<>c__DisplayClass4_0.<<ConfigureAwaitHideErrors>g__Inner|0>d.MoveNext() --- End of stack trace from previous location ---
当前代码实现
每5秒拉取消息的方法:
public async Task Invoke() { var subscriber = await SubscriberClient.CreateAsync(CreateSubscriptionName()); await subscriber.StartAsync((msg, cancellationToken) => { // Process the message... return Task.FromResult(SubscriberClient.Reply.Ack); }); await subscriber.StopAsync(CancellationToken.None); }
问题根源
官方明确提示:PublisherClient和SubscriberClient是成本高昂的对象,若需定期向同一主题发布或从同一订阅拉取消息,应创建单例实例并在应用生命周期内复用。
当前代码每5秒就创建一个新的SubscriberClient,用完即关闭,频繁的创建销毁会导致大量TCP连接被快速建立又释放,最终触发服务器端的连接重置(Connection reset by peer)。
解决方案
将SubscriberClient改为单例模式,在应用启动时创建一次,复用同一个实例处理所有拉取请求:
1. 实现单例客户端
public static class PubSubSubscriber { private static SubscriberClient _subscriber; private static readonly object _lock = new object(); public static async Task<SubscriberClient> GetSubscriberAsync() { if (_subscriber == null) { lock (_lock) { if (_subscriber == null) { _subscriber = await SubscriberClient.CreateAsync(CreateSubscriptionName()); // 启动一次订阅处理逻辑 await _subscriber.StartAsync((msg, cancellationToken) => { // Process the message... return Task.FromResult(SubscriberClient.Reply.Ack); }); } } } return _subscriber; } // 应用关闭时调用,清理资源 public static async Task StopAsync() { if (_subscriber != null) { await _subscriber.StopAsync(CancellationToken.None); _subscriber = null; } } private static SubscriptionName CreateSubscriptionName() { // 替换为你的项目ID和订阅ID return SubscriptionName.FromProjectSubscription("your-project-id", "your-subscription-id"); } }
2. 修改定时任务逻辑
定时任务不再创建新客户端,直接复用单例:
public async Task Invoke() { // 直接获取单例客户端,无需重复启动/停止 var subscriber = await PubSubSubscriber.GetSubscriberAsync(); // 客户端已在单例初始化时启动,此处无需额外操作 }
3. 应用生命周期管理
在应用启动时初始化单例,关闭时清理资源(以ASP.NET Core为例):
// Program.cs var builder = WebApplication.CreateBuilder(args); // 添加初始化服务 builder.Services.AddHostedService<PubSubInitializer>(); var app = builder.Build(); // ...其他配置 app.Run(); // 初始化服务类 public class PubSubInitializer : IHostedService { public async Task StartAsync(CancellationToken cancellationToken) { await PubSubSubscriber.GetSubscriberAsync(); } public async Task StopAsync(CancellationToken cancellationToken) { await PubSubSubscriber.StopAsync(); } }
额外注意事项
- 避免频繁启动/停止
SubscriberClient:客户端启动后会保持长连接持续拉取消息,无需每次定时任务都重启 - 保证消息处理逻辑线程安全:单例客户端会并发调用消息处理委托,需确保业务逻辑无线程安全隐患
- 优化拉取参数:可通过
SubscriberClient.Settings调整拉取批量大小、超时时间等参数,提升连接使用效率
内容的提问来源于stack exchange,提问作者jeff chen
相关产品推荐
相关产品推荐

