You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 13:50:31