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

基于.NET Core 3.1的C#程序迁移至Pulsar 2.10:返回消息处理问询

Pulsar 请求-回复模式的实现方案

Pulsar本身没有TIBCO那种原生的请求-回复API,但你可以通过以下几种手动实现的方式,复刻类似的同步/异步回复机制:

1. 专属回复主题的同步处理

这是最常用的方案,核心逻辑是为每个请求生成唯一标识,接收方处理完成后将回复发送到指定的专属回复主题,发送方订阅该主题并通过标识匹配回复:

C# 代码示例

  • 发送请求端:
var requestId = Guid.NewGuid().ToString();
// 为当前请求创建专属回复主题
var replyTopic = $"persistent://my-tenant/my-namespace/reply-{requestId}";

// 构造请求消息,携带请求ID和回复主题
var requestMsg = PulsarClient.NewMessage()
    .Topic("persistent://my-tenant/my-namespace/business-topic")
    .Property("requestId", requestId)
    .Property("replyTopic", replyTopic)
    .Content("业务请求内容")
    .Build();

await producer.SendAsync(requestMsg);

// 临时订阅专属回复主题,等待对应回复
using var replyConsumer = await pulsarClient.NewConsumer()
    .Topic(replyTopic)
    .SubscriptionName($"temp-sub-{requestId}")
    .SubscriptionType(SubscriptionType.Exclusive)
    .MessageFilter(msg => msg.Properties.TryGetValue("requestId", out var id) && id == requestId)
    .SubscribeAsync();

// 等待回复,设置超时避免无限阻塞
var replyMsg = await replyConsumer.ReceiveAsync(TimeSpan.FromSeconds(30));
var replyContent = Encoding.UTF8.GetString(replyMsg.Data);
// 处理回复逻辑
Console.WriteLine($"收到回复:{replyContent}");
await replyConsumer.AcknowledgeAsync(replyMsg);
  • 接收处理端:
var receivedMsg = await consumer.ReceiveAsync();
var requestId = receivedMsg.Properties["requestId"];
var replyTopic = receivedMsg.Properties["replyTopic"];

// 执行业务处理逻辑
var replyContent = "处理完成后的回复内容";

// 发送回复到指定主题,携带相同请求ID
var replyMsg = PulsarClient.NewMessage()
    .Topic(replyTopic)
    .Property("requestId", requestId)
    .Content(replyContent)
    .Build();

await replyProducer.SendAsync(replyMsg);
await consumer.AcknowledgeAsync(receivedMsg);

2. 共享回复主题+请求ID匹配

如果不想为每个请求创建独立主题,可以使用一个共享的回复主题,发送方通过请求ID过滤属于自己的回复:

  • 所有请求统一指定同一个共享回复主题
  • 发送方订阅该主题后,在接收消息时检查requestId是否匹配自己发送的请求标识
  • 注意:需要确保过滤逻辑严谨,避免接收其他请求的回复

3. 异步回调模式

如果不需要同步等待回复,可以采用异步回调的方式:

  • 维护一个线程安全的字典,将requestId映射到TaskCompletionSource或回调方法
  • 后台订阅回复主题,收到消息后根据requestId查找并触发对应的回调

C# 代码示例

private readonly ConcurrentDictionary<string, TaskCompletionSource<string>> _pendingRequests = new();

// 发送请求方法
public async Task<string> SendRequestAsync(string content)
{
    var requestId = Guid.NewGuid().ToString();
    var tcs = new TaskCompletionSource<string>();
    _pendingRequests.TryAdd(requestId, tcs);

    // 发送消息到业务主题,指定共享回复主题
    var requestMsg = PulsarClient.NewMessage()
        .Topic("business-topic")
        .Property("requestId", requestId)
        .Property("replyTopic", "shared-reply-topic")
        .Content(content)
        .Build();
    await producer.SendAsync(requestMsg);

    return await tcs.Task;
}

// 后台处理回复的方法
private async Task ListenForReplies()
{
    using var replyConsumer = await pulsarClient.NewConsumer()
        .Topic("shared-reply-topic")
        .SubscriptionName("reply-subscription")
        .SubscribeAsync();

    while (true)
    {
        var replyMsg = await replyConsumer.ReceiveAsync();
        if (replyMsg.Properties.TryGetValue("requestId", out var requestId) &&
            _pendingRequests.TryRemove(requestId, out var tcs))
        {
            var replyContent = Encoding.UTF8.GetString(replyMsg.Data);
            tcs.SetResult(replyContent);
        }
        await replyConsumer.AcknowledgeAsync(replyMsg);
    }
}

关键注意事项

  • 必须处理超时场景:如果在指定时间内未收到回复,需要取消等待并清理字典或临时订阅资源
  • 临时主题可设置自动过期策略,避免残留无用主题占用资源
  • 确保requestId的全局唯一性,防止回复匹配错误
  • 高并发场景下,使用线程安全集合(如ConcurrentDictionary)管理待处理请求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 08:13:14