【问题标题】:Adding Observers to already running MassTransit system将观察者添加到已经运行的 MassTransit 系统
【发布时间】:2018-06-06 00:32:56
【问题描述】:

我正在尝试将微服务添加到包含 MassTransit 观察者的系统中,该观察者将观察请求响应或发布系统中已使用的消息。我不能轻易地重新部署现有服务,所以如果可能的话,我宁愿避免它。

以下代码只在服务启动时执行,发送消息时不执行。

                BusControl = Bus.Factory.CreateUsingRabbitMq(cfg =>
                {
                    var host = cfg.Host(new Uri($"{settings.Protocol}://{settings.RabbitMqHost}/"), h =>
                    {
                        h.Username(settings.RabbitMqConsumerUser);
                        h.Password(settings.RabbitMqConsumerPassword);
                    });

                    cfg.ReceiveEndpoint(host, "pub_sub_flo", ec => { });

                    host.ConnectSendObserver(new RequestObserver());
                    host.ConnectPublishObserver(new RequestObserver());

                });

观察者:

 public class RequestObserver : ISendObserver, IPublishObserver
    {
        public Task PreSend<T>(SendContext<T> context) where T : class
        {
            return Task.CompletedTask;
        }

        public Task PostSend<T>(SendContext<T> context) where T : class
        {
            var proxy = new StoreProxyFactory().CreateProxy("fabric:/MessagePatterns");

            proxy.AddEvent(new ConsumerEvent()
            {
                Id = Guid.NewGuid(),
                ConsumerId = Guid.NewGuid(),
                Message = "AMQPRequestResponse",
                Date = DateTimeOffset.Now,
                Type = "Observer"
            }).Wait();

            return Task.CompletedTask;
        }

        public Task SendFault<T>(SendContext<T> context, Exception exception) where T : class
        {
            return Task.CompletedTask;
        }

        public Task PrePublish<T>(PublishContext<T> context) where T : class
        {
            return Task.CompletedTask;
        }

        public Task PostPublish<T>(PublishContext<T> context) where T : class
        {
            var proxy = new StoreProxyFactory().CreateProxy("fabric:/MessagePatterns");

            proxy.AddEvent(new ConsumerEvent()
            {
                Id = Guid.NewGuid(),
                ConsumerId = Guid.NewGuid(),
                Message = "AMQPRequestResponse",
                Date = DateTimeOffset.Now,
                Type = "Observer"
            }).Wait();

            return Task.CompletedTask;
        }

        public Task PublishFault<T>(PublishContext<T> context, Exception exception) where T : class
        {
            return Task.CompletedTask;
        }
    }

谁能帮忙?

非常感谢。

【问题讨论】:

    标签: masstransit


    【解决方案1】:

    仅在它们所连接的总线实例上发送、发布等消息时才会调用观察者。它们不会观察其他总线实例发送或接收的消息。

    如果您想观察这些消息,您可以创建一个观察者队列并将该队列绑定到您的服务交换,以便将请求消息的副本发送到您的服务。然而,回复并不容易获得,因为它们是通过临时交换直接发送到客户端队列的。

    cfg.ReceiveEndpoint(host, "service-observer", e =>
    {
        e.Consumer<SomeConsumer>(...);
        e.Bind("service-endpoint");
    });
    

    这会将服务端点交换绑定到您的接收端点队列,以便将消息的副本发送给您的消费者。

    这通常被称为丝锥。

    【讨论】:

    • 谢谢,这正是我遇到的问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-11-18
    • 2022-08-12
    • 2012-01-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-05-24
    相关资源
    最近更新 更多