代码之家  ›  专栏  ›  技术社区  ›  RoLYroLLs

具有工作角色的Azure服务总线主题和订阅

  •  0
  • RoLYroLLs  · 技术社区  · 7 年前

    所以我最近需要用 Service Bus Topic and Subscriptions 我也看过很多文章和教程。我已经成功地实现了微软的 Get started with Service Bus topics 也被成功使用 Visual Studio 2017's Worker Role 访问数据库的模板。

    然而,我对如何将两者恰当地“结合”起来感到困惑。而 服务总线主题入门 文章展示了如何创建两个应用程序,一个用于发送,一个用于接收,然后退出 工作者角色 模板似乎无休止地循环 await Task.Delay(10000); .

    我不知道怎么把这两个“啮合”好。本质上,我想 工作者角色 为了保持活力,永远聆听订阅中的条目(或者直到它明显退出)。

    任何指导都很好!

    P、 S:我问了一个相关的问题,关于我应该在我的案例场景中使用的适当技术 StackExchange - Software Engineering 如果你感兴趣的话。

    更新#1(2018/08/09)

    基于 Arunprabhu's answer ,以下是我发送 Message 根据我阅读和接收的文章 Visual Studio 2017 Worker Role with Service Bus Queue 模板。

    发送(基于 服务总线主题入门 )

    using System;
    using System.Text;
    using System.Threading.Tasks;
    using Microsoft.Azure.ServiceBus;
    
    namespace TopicsSender {
        internal static class Program {
            private const string ServiceBusConnectionString = "<your_connection_string>";
            private const string TopicName = "test-topic";
            private static ITopicClient _topicClient;
    
            private static void Main(string[] args) {
                MainAsync().GetAwaiter().GetResult();
            }
    
            private static async Task MainAsync() {
                const int numberOfMessages = 10;
                _topicClient = new TopicClient(ServiceBusConnectionString, TopicName);
    
                Console.WriteLine("======================================================");
                Console.WriteLine("Press ENTER key to exit after sending all the messages.");
                Console.WriteLine("======================================================");
    
                // Send messages.
                await SendMessagesAsync(numberOfMessages);
    
                Console.ReadKey();
    
                await _topicClient.CloseAsync();
            }
    
            private static async Task SendMessagesAsync(int numberOfMessagesToSend) {
                try {
                    for (var i = 0; i < numberOfMessagesToSend; i++) {
                        // Create a new message to send to the topic
                        var messageBody = $"Message {i}";
                        var message = new Message(Encoding.UTF8.GetBytes(messageBody));
    
                        // Write the body of the message to the console
                        Console.WriteLine($"Sending message: {messageBody}");
    
                        // Send the message to the topic
                        await _topicClient.SendAsync(message);
                    }
                } catch (Exception exception) {
                    Console.WriteLine($"{DateTime.Now} :: Exception: {exception.Message}");
                }
            }
        }
    }
    

    接收(基于 具有服务总线队列的工作角色 模板)

    using System;
    using System.Diagnostics;
    using System.Net;
    using System.Threading;
    using Microsoft.ServiceBus.Messaging;
    using Microsoft.WindowsAzure.ServiceRuntime;
    
    namespace WorkerRoleWithSBQueue1 {
        public class WorkerRole : RoleEntryPoint {
            // The name of your queue
            private const string ServiceBusConnectionString = "<your_connection_string>";
            private const string TopicName = "test-topic";
            private const string SubscriptionName = "test-sub1";
    
            // QueueClient is thread-safe. Recommended that you cache 
            // rather than recreating it on every request
            private SubscriptionClient _client;
            private readonly ManualResetEvent _completedEvent = new ManualResetEvent(false);
    
            public override void Run() {
                Trace.WriteLine("Starting processing of messages");
    
                // Initiates the message pump and callback is invoked for each message that is received, calling close on the client will stop the pump.
                _client.OnMessage((receivedMessage) => {
                    try {
                        // Process the message
                        Trace.WriteLine("Processing Service Bus message: " + receivedMessage.SequenceNumber.ToString());
                        var message = receivedMessage.GetBody<byte[]>();
                        Trace.WriteLine($"Received message: SequenceNumber:{receivedMessage.SequenceNumber} Body:{message.ToString()}");
                    } catch (Exception e) {
                        // Handle any message processing specific exceptions here
                        Trace.Write(e.ToString());
                    }
                });
    
                _completedEvent.WaitOne();
            }
    
            public override bool OnStart() {
                // Set the maximum number of concurrent connections 
                ServicePointManager.DefaultConnectionLimit = 12;
    
                // Initialize the connection to Service Bus Queue
                _client = SubscriptionClient.CreateFromConnectionString(ServiceBusConnectionString, TopicName, SubscriptionName);
                return base.OnStart();
            }
    
            public override void OnStop() {
                // Close the connection to Service Bus Queue
                _client.Close();
                _completedEvent.Set();
                base.OnStop();
            }
        }
    }
    

    更新2(2018/08/10)

    在收到Arunprabhu的一些建议并知道我正在使用不同的库之后,下面是我当前的解决方案,其中的部分来自多个来源。有没有什么我忽略了,加上肩在那里,等等?目前得到一个错误,可能是为另一个问题或已经回答,所以不想张贴它之前,进一步的研究。

    using System;
    using System.Diagnostics;
    using System.Net;
    using System.Text;
    using System.Threading;
    using System.Threading.Tasks;
    using Microsoft.Azure.ServiceBus;
    using Microsoft.WindowsAzure.ServiceRuntime;
    
    namespace WorkerRoleWithSBQueue1 {
        public class WorkerRole : RoleEntryPoint {
            private readonly CancellationTokenSource _cancellationTokenSource = new CancellationTokenSource();
            private readonly ManualResetEvent _runCompleteEvent = new ManualResetEvent(false);
    
            // The name of your queue
            private const string ServiceBusConnectionString = "<your_connection_string>";
            private const string TopicName = "test-topic";
            private const string SubscriptionName = "test-sub1";
    
            // _client is thread-safe. Recommended that you cache 
            // rather than recreating it on every request
            private SubscriptionClient _client;
    
            public override void Run() {
                Trace.WriteLine("Starting processing of messages");
    
                try {
                    this.RunAsync(this._cancellationTokenSource.Token).Wait();
                } catch (Exception e) {
                    Trace.WriteLine("Exception");
                    Trace.WriteLine(e.ToString());
                } finally {
                    Trace.WriteLine("Finally...");
                    this._runCompleteEvent.Set();
                }
            }
    
            public override bool OnStart() {
                // Set the maximum number of concurrent connections 
                ServicePointManager.DefaultConnectionLimit = 12;
    
                var result = base.OnStart();
    
                Trace.WriteLine("WorkerRole has been started");
    
                return result;
            }
    
            public override void OnStop() {
                // Close the connection to Service Bus Queue
                this._cancellationTokenSource.Cancel();
                this._runCompleteEvent.WaitOne();
    
                base.OnStop();
            }
    
            private async Task RunAsync(CancellationToken cancellationToken) {
                // Configure the client
                RegisterOnMessageHandlerAndReceiveMessages(ServiceBusConnectionString, TopicName, SubscriptionName);
    
                _runCompleteEvent.WaitOne();
    
                Trace.WriteLine("Closing");
                await _client.CloseAsync();
            }
    
            private void RegisterOnMessageHandlerAndReceiveMessages(string connectionString, string topicName, string subscriptionName) {
                _client = new SubscriptionClient(connectionString, topicName, subscriptionName);
    
                var messageHandlerOptions = new MessageHandlerOptions(ExceptionReceivedHandler) {
                    // Maximum number of concurrent calls to the callback ProcessMessagesAsync(), set to 1 for simplicity.
                    // Set it according to how many messages the application wants to process in parallel.
                    MaxConcurrentCalls = 1,
    
                    // Indicates whether MessagePump should automatically complete the messages after returning from User Callback.
                    // False below indicates the Complete will be handled by the User Callback as in `ProcessMessagesAsync` below.
                    AutoComplete = false,
                };
    
                _client.RegisterMessageHandler(ProcessMessageAsync, messageHandlerOptions);
            }
    
            private async Task ProcessMessageAsync(Message message, CancellationToken token) {
                try {
                    // Process the message
                    Trace.WriteLine($"Received message: SequenceNumber:{message.SystemProperties.SequenceNumber} Body:{Encoding.UTF8.GetString(message.Body)}");
                    await _client.CompleteAsync(message.SystemProperties.LockToken);
                } catch (Exception e) {
                    // Handle any message processing specific exceptions here
                    Trace.Write(e.ToString());
                    await _client.AbandonAsync(message.SystemProperties.LockToken);
                }
            }
    
            private static Task ExceptionReceivedHandler(ExceptionReceivedEventArgs exceptionReceivedEventArgs) {
                Console.WriteLine($"Message handler encountered an exception {exceptionReceivedEventArgs.Exception}.");
                var context = exceptionReceivedEventArgs.ExceptionReceivedContext;
                Console.WriteLine("Exception context for troubleshooting:");
                Console.WriteLine($"- Endpoint: {context.Endpoint}");
                Console.WriteLine($"- Entity Path: {context.EntityPath}");
                Console.WriteLine($"- Executing Action: {context.Action}");
                return Task.CompletedTask;
            }
        }
    }
    
    3 回复  |  直到 7 年前
        1
  •  2
  •   Arunprabhu    7 年前

    考虑到更新问题更新1(2018/08/09)的复杂性,我提供了一个单独的答案。

    发送方和接收方正在使用不同的库。

    发送者- Microsoft.Azure.ServiceBus

    接收器- WindowsAzure.ServiceBus

    Microsoft.Azure.ServiceBus的消息对象为message,其中WindowsAzure.ServiceBus已代理消息。

    有一种方法 RegisterMessageHandler 在Microsoft.Azure.ServiceBus中可用,这是WindowsAzure.ServiceBus中client.OnMessage()的替代方法。通过使用它,侦听器将消息作为消息对象接收。这个库支持您期望的异步编程。

    参考 here 两个库的样本。

        2
  •  1
  •   Arunprabhu    7 年前

    如果您使用的是Visual Studio,则有一个默认模板可用于创建带有服务总线队列的Azure云服务和工作者角色。您需要在WorkerRole.cs中使用SubscriptionClient更改QueueClient。

    然后,worker角色将保持活动状态,侦听来自主题订阅的消息。

    你可以找到样品 here . 您应该在云服务中使用服务总线队列创建工作角色

    enter image description here

        3
  •  0
  •   Nyamnyam    7 年前

    我想知道你是否有任何特定的理由选择员工角色而不是网络工作?如果没有,您可以使用具有ServiceBusTrigger属性的Web作业,它将使您的代码更加简单。 more info here...