基于.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
相关产品推荐
相关产品推荐

