在本视频中,我们将探讨…… Pub/Sub 消息传递 - NCache 使用 .NET 应用程序。发布/订阅是一种消息传递模式,它允许系统的不同部分相互通信,而无需直接依赖彼此。因此,系统中有一个发布者发送消息,一个订阅者监听消息。中间有一个代理,用于将消息从一端路由到另一端。
这是哪里 NCache 它进来了。它在这里充当经纪人,使用这个主题。正如你在这里看到的。现在这个主题位于…… NCache 它缓存集群并存储您的消息。它存储您的发布者信息和订阅者信息,并且还用于将事件转发给您的发布者或订阅者。
这一切,就是为了让你的消息能够从一端发送到另一端,而无需应用程序之间相互依赖。
首先,我们来看看如何创建一个主题。正如你所看到的,我们在缓存中调用了消息服务接口,也就是创建主题的方法。当然,我们唯一需要传入的参数就是主题名称。主题创建完成后,它会一直存在于缓存中,直到你将其删除。当然,上面这行代码就是用来获取主题的。创建主题后,你只需要调用这个方法即可。
创建/获取主题
String topicName = "ExampleTopic";
ITopic topic = _cache.MessagingService.GetTopic(topicName);
if (topic == null)
{
topic = _cache.MessagingService.CreateTopic(topicName);
}
交货期权
投放方式
我们已经创建了一个主题。接下来我们要做的就是向该主题发布一条消息。 NCache 它提供了几种不同的消息发布功能。您可以选择不同的发送方式。您可以选择“全部”将消息发送给所有注册订阅者,也可以选择“任意”仅向其中一位订阅者发送消息。
此外,您可以选择消息发送模式,以便同步发布消息。您也可以异步发送消息,这当然能提供更好的性能。最后,您还可以批量发送消息,这样会将消息分批发送,而不是逐条发送。批量发送当然也能提供更好的性能。如果您希望订阅者应用程序按照发布者发送消息的顺序接收消息,则应选择同步发送。这就是您选择同步发送的原因。
现在我们来看一下发布消息的代码。当然,首先我们需要获取 ITopic 对象,也就是我们的主题对象。我们创建一个消息对象,传入的参数是有效负载。我这里传入的是一个订单对象作为有效负载。你还可以为消息设置过期时间。最后,我们调用 publish 方法。我们传入这些参数:消息本身、传递方式,你还可以异步发布消息。你可以看到,这里的第三个参数是用于消息传递失败通知的。我稍后会解释一下。
ITopic topic = cache.MessagingService.GetTopic(topicName);
Message message = new Message(new Order(), TimeSpan.FromSeconds(15));
// deliver message to all subscribers
topic.Publish(message, DeliveryOption.All, true);
// deliver message to any one subscriber
topic.Publish(message, DeliveryOption.Any, true);
// send message asynchronously
topic.PublishAsync(message, deliveryOption, true);
所以,你需要启用那个选项并将其设置为 true,同时还要创建这个回调方法。在这个例子中,你可以看到它只是输出一条日志消息,提示消息传递失败。这行代码是用来注册事件处理程序的,所以你需要在上一张幻灯片中启用那个参数,你可以在这里看到。然后,你可以通过这段代码注册事件处理程序。这样,每次订阅者没有收到你的消息时,发布者都会收到通知。之后,你可以使用这个回调函数,让你的应用程序做出相应的反应,或者记录一条消息。
private static void MessageDeliveryFailureNotification(object sender, MessageFailedEventArgs args) {
Console.WriteLine("Failed to deliver message. " + args.MessageFailureReason);
}
// register event handler
ITopic topic = cache.MessagingService.GetTopic(topicName);
topic.MessageDeliveryFailure += MessageDeliveryFailedNotification;
订阅类型
方针政策
投放方式
所以你可以看到,可以有多个发布者向某个特定主题发布消息。在该主题下,可以有多个订阅,当然,每个订阅又可以有多个订阅者。你可以看到,我们实际上可以有不同类型的订阅。有持久订阅,也有非持久订阅。
持久订阅的好处在于,您可以将其视为永久订阅。即使您是订阅用户,但您的应用程序离线或断开连接,它也会保存您的消息。消息将一直保存到重新连接或消息过期为止。
与此相反,对于非持久订阅,如果您的订阅应用程序断开连接,消息将会丢失。因此,请将它们视为临时订阅。
然后您还可以选择订阅策略。“独享”策略最多只允许一位订阅者。“共享”策略允许多位订阅者。
所以,非长期订阅默认是独占的,最多只能有一个订阅者;而长期订阅当然可以自定义订阅策略。好的。
最后,与发布者的工作方式类似,订阅也可以设置消息发送模式为同步或异步。当然,其中一种模式侧重于消息排序,而另一种模式则提供更佳的性能。
我们先来看看如何创建一个非持久订阅。可以看到,我们首先创建一个主题,然后调用创建订阅的方法。这样就创建了一个非持久订阅。这里我们传递的唯一方法参数是消息接收回调函数。好的。消息接收回调函数的作用是,当订阅者收到消息时通知它。在这个例子中,它只是简单地输出一条消息。
ITopic topic = cache.MessagingService.GetTopic(topicName);
// Create and register subscribers for the given topic
// Message received callback has to be passed
ITopicSubscription subscription = topic.CreateSubscription(MessageReceivedCallback);
// Used to notify the Subscriber when it receives a message
private static void MessageReceivedCallback(object sender, MessageEventArgs args) {
Console.WriteLine("Message received for topic " + args.TopicName);
}
当然,持久订阅的情况略有不同,你需要传入多个参数。需要特别说明的是,持久订阅是命名订阅,你可以给它命名。而非持久订阅则不能命名。所以,创建持久订阅时,我们传入的第一个参数是名称,然后是策略(共享或独占),消息接收回调以及订阅的到期时间。好的。
// Multiple Subscribers can subscribe to this Subscription
SubscriptionPolicy sharedPolicy = SubscriptionPolicy.Shared;
// Only one Subscriber allowed on this Subscription
SubscriptionPolicy exclusivePolicy = SubscriptionPolicy.Exclusive;
IDurableTopicSubscription subscription =
topic.CreateDurableSubscription("ExampleSubscription",
sharedPolicy,
MessageReceivedCallback,
TimeSpan.FromHours(1));
那么,在进入下一张幻灯片之前,我想先展示一下我的示例应用程序。当然,这是一个 .NET 应用程序。您可以看到我已经编写好了代码。这里有四个解决方案。当然,这里还有示例数据。这是一个订单。我们有订单发布者、主订单订阅者和辅助订单订阅者。好的。
所以你可以看到,在我的订单发布器中,我有两个方法。这个方法是异步发布订单。当然,我们首先创建主题对象。注册消息传递失败事件,然后遍历所有订单并将它们作为消息发布。这条消息的过期时间设置为 15 秒。
我的发送选项,在本例中设置为所有订阅者。好的。每条消息发送完毕后,您可以看到我为了演示目的添加了五秒的延迟。它应该会在发送后输出每个订单。这里还有异步选项。它会调用异步发布方法。然后我们还有批量发布订单的功能。好的。这就是我的消息发送失败通知的样子。您可以看到它输出了这条日志消息“发送订单 ID 失败”,然后也发送了订单详情。好的。接下来,我将在不运行订阅者的情况下运行此应用程序。
但在那之前,让我们先来看看我们的 NCache 缓存集群。我已经创建了一个名为 demoCache 的缓存。这是一个单服务器节点,只有一个服务器连接到它,并且部署在 Docker 实例上。我只想检查一下统计数据,确保这里没有数据。完美。好的。缓存已启动并运行。我现在可以运行我的应用程序了。我只是想尽量最小化进程,以便并排显示统计数据。
好了,我们来运行发布者。我先解释一下我预期看到的情况:由于目前还没有订阅者应用程序在运行,它会在缓存中创建一个主题。首先,我们可以看到客户端已连接。它创建了一个主题,现在正在同步地逐个发布订单。可以看到缓存大小也在增加,因为主题已存储在缓存中。
由于没有订阅者应用程序监听这些订单,我希望发布者收到消息发送失败的通知。可以看到,订单 ID 为 1 的消息发送失败,原因是消息已过期。我希望所有五个订单都出现同样的情况。所以,是的,订单发送失败了。为了节省大家的时间,我将关闭此问题。
接下来我想演示的是先运行订阅者应用程序,然后再运行发布者应用程序。好的。在此之前,让我先给你们看看这边的代码。可以看到,我首先在这里初始化了缓存。我们获取了订单主题。我们创建了一个名为“订单订阅”的持久订阅。但这只是我的私有方法,并不是实际的 API 调用。
这就是 API 调用的样子。好的。我们传入了订阅名称。策略,我认为在这个例子中是共享策略。所以你可以有多个订阅者。我们还传入了消息接收回调函数和一个小时的时间跨度。回调函数很简单,它输出“已接收订单”,并打印订单详情,就像你在这里看到的那样。现在我要运行这个应用程序,它应该会创建我的订阅。好的。创建完成后,我将运行我的发布者应用程序,看看它的表现如何。
这是我的订阅者,这是我的发布者。订单号为 1 的消息已经发布,您可以看到我们的订阅者正在接收所有订单。您会注意到,由于我们使用的是同步发布,而我们的订阅者应用程序也默认同步处理,所以您可以看到它按顺序接收了所有订单。我认为订单号为 5 的消息发送失败是上一条消息中的错误,因为我在发送之前关闭了应用程序,所以我们暂时忽略它。但您可以看到一切都井然有序。
现在我们来尝试异步发布消息,看看它的行为有何不同。我把这个方法调用改成异步的,参数应该保持不变。我把这里的过期时间改成 60 秒,这样我们就可以忽略缓冲区延迟。接下来我要展示的是异步和同步的区别。我们运行订阅者应用程序,确保所有程序都已关闭。好的。订阅者应用程序启动运行后,我们运行发布者应用程序。我预期现在看到的是订阅者应用程序接收到所有订单,但订单的接收顺序并不固定。可以看到,所有订单都已发布,我们也确实收到了订单,但顺序是这样的:先收到订单号 3,然后收到订单号 2,接着收到订单号 5,最后才收到订单号 5 和订单号 1。所以订单的接收顺序并不固定,这就是异步和同步的区别。
当然,在订阅者这边,我简单演示一下。我们选择了同步模式进行订单处理,但为了确保它100%正常工作,您需要在订阅者和发布者之间都进行同步设置。您可以看到,这里默认选择的是同步交付模式。您需要通过一个额外的参数来更改同步和异步模式。好的。
所以我们来试试别的。我们要看看持久订阅和非持久订阅的区别。首先,我来演示一下持久订阅如何在应用程序离线的情况下保留消息。好的。首先,我确保所有程序都已关闭。我们要创建一个持久订阅。我们先运行一下这个应用程序看看。现在,为了演示,我将发布方式改回同步发布。请原谅我的拼写错误。我们将同步发布订单,并将过期时间改为 60 秒,以避免缓冲区延迟。现在运行一下。在发布者发布消息的同时,我将中途关闭订阅者应用程序,然后再重新启动,看看它的表现如何。
现在它们都在运行。第一个订单已收到,我要把它关闭。它会继续向主题发布所有订单。现在我要恢复订阅者。请注意,它实际上不会创建一个新的订阅。因为它是持久的,所以它只会恢复当前存在的订阅。可以看到,在我的应用程序离线期间发布了几个订单,它仍然能够接收它们。我们已经收到了第五个订单,我预计它还会收到第三个和第四个订单,也许还有第二个。所以第四个订单也已收到。我们快速并排打开缓存统计信息。可以看到有两个客户端连接。一个是订阅者应用程序,一个是发布者应用程序。希望我们能够在消息过期之前收到它们,这就是我设置 60 秒时间的原因。所以,现在应该就在几秒钟之内了。当然,希望消息不会过期。这是第三个订单。好的。完美。
我现在要结束这个演示,接下来我们要看看在完全相同的情况下,非持久订阅的行为如何。首先,我来展示一下非持久订阅的代码。可以看到,我们获取了订单主题。然后,我们调用了创建订阅的方法,并传入了回调函数。这样就创建了一个非持久订阅。
ITopic topic = cache.MessagingService.GetTopic(topicName);
// Create and register subscribers for the given topic
// Message received callback has to be passed
ITopicSubscription subscription = topic.CreateSubscription(MessageReceivedCallback);
// Used to notify the Subscriber when it receives a message
private static void MessageReceivedCallback(object sender, MessageEventArgs args) {
Console.WriteLine("Message received for topic " + args.TopicName);
}
我现在就运行这个程序。运行之后,我会启动发布者应用程序。现在它已经启动并运行了,可以接收订单了。我先把它关闭,等几个订单发布后再重新运行订阅者应用程序。现在我们重新运行一下。可以看到,订单三和订单四是在应用程序离线时发布的,所以它无法接收这些订单。订单五是在应用程序运行期间接收到的,所以你可以看到它显示出来。但显然它无法接收其他订单,因为这些消息丢失了。为了演示,我再等大约十秒钟,向你展示它无法保留消息。
所以,在进行上述操作的同时,我也想向您展示如何确保消息的有序性。正如我之前解释过的,您需要确保发布者同步发送消息,然后您还可以使用订阅者同步处理消息。 NCache 除此之外,它还提供了另一个功能。您可以在发布方法调用中指定一个序列名称,这样就能确保特定批次的消息按顺序发送。好的。在这段代码中,您可以看到,我们创建了一条消息,并使用这个序列名称发布了该消息。我们对整批消息都执行此操作。因此,这当然也能确保消息的有序性。
ITopic topic = cache.MessagingService.GetTopic(topicName);
for (int i = 0; i < 30; i++) {
Order order = FetchAnyOrderFromDB();
Message message = new Message(order);
// Specify a unique sequence name for the messages
string sequenceName = "OrderMessages";
// Publish message with the sequence name
topic.Publish(message, DeliveryOption.All, sequenceName, true);
}
至于订阅者方面,我之前也演示过了,你需要传递传递模式,这是一个额外的参数,用于选择同步或异步传递模式。好的。
ITopic topic = cache.MessagingService.GetTopic(topicName);
// Create and register subscribers for Topic
// Message received callback is specified
// DeliveryMode is set to async to ensure ordered messages
ITopicSubscription subscription = topic.CreateSubscription(MessageReceived, DeliveryMode.Sync);
回到订单输出,可以看到消息已经过期,无法接收。所以一切基本都按预期运行。
各位,演示就到这里。希望你们喜欢这个视频。如果想了解更多信息,请继续阅读。 NCache,它的 功能如果您想了解如何将其融入您的申请中,请联系我们。 安排演示非常感谢您观看本视频。