BestQA.EventSubscribe/ENodeExtensions.cs

(开头部分) 2KB

这里只显示每个文件的开头 60 行。登录后可以解锁完整代码。

using System.Linq;
using System.Threading;
using BestQA.RegisterExtension;
using ECommon.Components;
using ECommon.Logging;
using ECommon.Scheduling;
using ENode.Commanding;
using ENode.Configurations;
using ENode.EQueue;
using ENode.Infrastructure;
using ENode.Infrastructure.Impl;
using EQueue.Configurations;
namespace BestQA.EventSubscribe
{
    public static class ENodeExtensions
    {
        private static DomainEventConsumer _eventConsumer;
        private static CommandService _commandService;
        public static ENodeConfiguration UseEQueue(this ENodeConfiguration enodeConfiguration)
        {
            var configuration = enodeConfiguration.GetCommonConfiguration();
            configuration.RegisterEQueueComponents();
            _commandService = new CommandService(id: "CommandServiceForProcessManager");
            configuration.SetDefault<ICommandService, CommandService>(_commandService);
            _eventConsumer = new DomainEventConsumer();
            _eventConsumer.Subscribe("QuestionEventTopic");
            return enodeConfiguration;
        }
        public static ENodeConfiguration StartEQueue(this ENodeConfiguration enodeConfiguration)
        {
            _commandService.Start();
            _eventConsumer.Start();
            WaitAllConsumerLoadBalanceComplete();
            return enodeConfiguration;
        }
        public static ENodeConfiguration RegisterAllTypeCodes(this ENodeConfiguration enodeConfiguration)
        {
            var provider = ObjectContainer.Resolve<ITypeCodeProvider>() as DefaultTypeCodeProvider;
            provider.RegisterEvent().RegisterEventHandler();
            return enodeConfiguration;
        }
        private static void WaitAllConsumerLoadBalanceComplete()
        {
            var logger = ObjectContainer.Resolve<ILoggerFactory>().Create(typeof(ENodeExtensions).Name);
            var scheduleService = ObjectContainer.Resolve<IScheduleService>();
            var waitHandle = new ManualResetEvent(false);
            logger.Info("等待所有负载加载...");
            var taskId = scheduleService.ScheduleTask("WaitAllConsumerLoadBalanceComplete", () =>
            {
                var eventConsumerAllocatedQueues = _eventConsumer.Consumer.GetCurrentQueues();
                if (eventConsumerAllocatedQueues.Count() == 8)
                {
                    waitHandle.Set();
                }
            }, 1000, 1000);
            waitHandle.WaitOne();
            scheduleService.ShutdownTask(taskId);
            logger.Info("All consumer load balance completed.");
        }
    }
后面还有 1 行代码,解锁后查看完整代码

24 小时内免费解锁 3 个项目,之后 1 积分/个。 规则说明

AI 解读

登录后可用,每次 10 积分,解读结果公开显示在下面。

还没有人解读过这个文件。