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

C#监听CouchDB _changes的continuous模式时程序挂起问题求解

问题根因

代码切换到continuous模式挂死的核心原因有两点:

  • .Result同步阻塞异步调用本身存在死锁风险,更关键的是continuous模式下CouchDB _changes接口不会返回带结束标识的完整响应:连接建立后服务端会持续保持连接,每产生一条变更就向响应流写入一行独立JSON,直到连接主动断开。ReadAsStringAsync()会等待整个响应内容接收完成才返回,因此会无限阻塞。
  • 原有逻辑等待完整响应后一次性反序列化的处理方式,完全不匹配continuous模式逐行输出JSON片段的响应格式。
实现方案

核心思路是跳过完整响应等待环节,直接以流模式逐行读取响应内容,每读到一条完整变更就立刻处理,同时补充重连逻辑保证连接可靠性。

前置配置

初始化HttpClient时关闭默认请求超时,避免长连接被客户端主动断开:

_httpClient.Timeout = Timeout.InfiniteTimeSpan;

核心监听代码

internal async Task ListenAsync(JobList jobList, CancellationToken cancellationToken)
{
    // 记录最后处理的变更序列号,重连时从断点续传
    long lastSequence = 0;
    while (!cancellationToken.IsCancellationRequested)
    {
        try
        {
            // 构造请求地址,每次重连带上lastSequence避免漏数/重复处理
            var requestUri = $"/你的数据库名/_changes?filter=_doc_ids&feed=continuous&include_docs=true&heartbeat=5000&since={lastSequence}";
            using var request = new HttpRequestMessage(HttpMethod.Get, requestUri);
            
            // 关键参数:拿到响应头立刻返回,不等待完整响应内容加载
            using var response = await _httpClient.SendAsync(
                request, 
                HttpCompletionOption.ResponseHeadersRead, 
                cancellationToken);
            response.EnsureSuccessStatusCode();

            // 直接获取响应流逐行读取
            await using var responseStream = await response.Content.ReadAsStreamAsync(cancellationToken);
            using var streamReader = new StreamReader(responseStream);

            while (!streamReader.EndOfStream && !cancellationToken.IsCancellationRequested)
            {
                var line = await streamReader.ReadLineAsync(cancellationToken);
                // 心跳包是空行,直接跳过
                if (string.IsNullOrWhiteSpace(line)) continue;

                // 单条变更单独反序列化
                var changeItem = JsonConvert.DeserializeObject<ChangeEntry>(line);
                if (changeItem == null) continue;

                // 更新断点序列号
                lastSequence = changeItem.Seq;

                // 处理变更内容,写入任务列表
                if (changeItem.Doc?.ImportJobs != null)
                {
                    foreach (var job in changeItem.Doc.ImportJobs)
                    {
                        jobList.addJob(job);
                    }
                }
            }
        }
        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
        {
            // 收到取消信号,正常退出监听
            break;
        }
        catch (Exception ex)
        {
            // 此处记录错误日志
            // 连接异常后等待3秒自动重连
            await Task.Delay(3000, cancellationToken);
        }
    }
}

// 反序列化用的实体类,根据实际返回结构调整即可
public class ChangeEntry
{
    [JsonProperty("seq")]
    public long Seq { get; set; }
    [JsonProperty("id")]
    public string DocId { get; set; }
    [JsonProperty("doc")]
    public JobDoc Doc { get; set; }
}

public class JobDoc
{
    [JsonProperty("ImportJobs")]
    public List<ImportJob> ImportJobs { get; set; }
}

启动监听时直接用异步任务启动即可,不需要额外开同步线程阻塞:

var cts = new CancellationTokenSource();
// 启动监听任务
_ = ListenAsync(jobList, cts.Token);

// 需要停止监听时调用
// cts.Cancel();
关键注意事项
  • SendAsync必须传入HttpCompletionOption.ResponseHeadersRead,这是流模式读取长连接的核心配置,缺省参数下HttpClient会缓存整个响应直到结束,依然会造成阻塞。
  • 全程使用await异步等待,不要用.Result、.Wait()这类同步阻塞方法,避免同步上下文死锁。
  • 必须维护lastSequence断点值,重连时通过since参数从断点续传,避免连接闪断时漏掉变更或者重复处理已处理内容。
  • 若使用_doc_ids过滤器,CouchDB要求通过POST请求将待监听的文档ID列表放在请求体中传递,GET请求携带该过滤器不会生效,需要核对请求构造逻辑是否符合要求。
  • 心跳逻辑不需要额外处理,读到空行直接跳过即可,该机制用于避免中间网络设备因为连接空闲切断长连接。

内容的提问来源于stack exchange,提问作者Kim Sandberg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 23:12:33