【问题标题】:MassTransit with AWS SQS - exception sending dynamic message带有 AWS SQS 的 MassTransit - 发送动态消息的异常
【发布时间】:2020-12-02 10:53:51
【问题描述】:

我将 Masstransit 与 AWS SQS/SNS 一起用作传输。 现在我遇到了一个简单的问题 - 将 CorrelationId 添加到已发布的消息中。由于我遵循使用接口作为 DTO 的建议,所以我不能使用CorrelatedBy<Guid> 接口,我决定直接将__CorrelationId 属性添加到消息中。

我正在使用以下包装器发布消息:

public class MessageBus : IMessageBus
{
    // ... omitted for brevity

    public async Task Publish<T>(object message)
        where T : class
    {
        await this.publishEndpoint.Publish<T>(this.CreateCorrelatedMessage(message));
    }

    private object CreateCorrelatedMessage(object message)
    {
        dynamic corellatedMessage = new ExpandoObject();

        foreach (var property in message.GetType().GetProperties())
        {
            ((IDictionary<string, object>)corellatedMessage).Add(property.Name, property.GetValue(message));
        }

        corellatedMessage.__CorrelationId = this.correlationIdProvider.CorrelationId;
        return corellatedMessage;
    }
}

这是用法:

        await this.messageBus.Publish<ObtainRequestToken>(new
        {
            Social = SocialEnum.Twitter,
            AppId = credentials.First(x => x.Name == "TwitterConsumerKey").Value,
            AppSecret = credentials.First(x => x.Name == "TwitterConsumerSecret").Value,
            ReturnUrl = returnUrl
        });

这在具有以下总线注册的控制台应用程序中运行良好:

        var busControl = Bus.Factory.CreateUsingAmazonSqs(cfg =>
        {
            cfg.Host("eu-west-1", h =>
            {
                h.AccessKey("xxx");
                h.SecretKey("xxx");

                h.Scope($"Local", scopeTopics: true);
            });
        });

但是 ASP.NET Core 应用程序与

        services.AddMassTransit(cfg =>
        {
            cfg.UsingAmazonSqs((context, sqsCfg) =>
            {
                var amazonAccount = this.Configuration
                    .GetSection("AmazonAccount")
                    .Get<AmazonAccountConfig>();

                sqsCfg.Host(amazonAccount.Region, h =>
                {
                    h.AccessKey(amazonAccount.KeyId);
                    h.SecretKey(amazonAccount.SecretKey);

                    h.Scope(this.Environment.EnvironmentName, scopeTopics: true);
                });
            });

            cfg.SetEndpointNameFormatter(new DefaultEndpointNameFormatter(this.Environment.EnvironmentName + "_", false));
        });

        services.AddMassTransitHostedService();

使用Object reference not set to an instance of an object. 和堆栈跟踪发布失败

   at MassTransit.Logging.EnabledDiagnosticSource.StartSendActivity[T](SendContext1 context, ValueTuple2[] tags)
   at MassTransit.AmazonSqsTransport.Transport.TopicSendTransport.SendPipe1.<Send>d__5.MoveNext()
   at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
   at System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
   at System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
   at System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
   at GreenPipes.Agents.PipeContextSupervisor1.<GreenPipes-IPipeContextSource<TContext>-Send>d__7.MoveNext()
   at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
   at GreenPipes.Agents.PipeContextSupervisor1.<GreenPipes-IPipeContextSource<TContext>-Send>d__7.MoveNext()
   at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
   at GreenPipes.Agents.PipeContextSupervisor1.<GreenPipes-IPipeContextSource<TContext>-Send>d__7.MoveNext()
   at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
   at System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
   at System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
   at System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
   at MassTransit.Initializers.MessageInitializer2.<Send>d__9.MoveNext()
   at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
   at System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
   at System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
   at System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
   at MassTransit.Transports.PublishEndpoint.<>c__DisplayClass18_01.<<PublishInternal>g__PublishAsync|0>d.MoveNext()
   at System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
   at System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
   at System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
   at System.Dynamic.UpdateDelegates.UpdateAndExecuteVoid1[T0](CallSite site, T0 arg0)
   at xxx.Main.Api.Messaging.MessageBus.<Publish>d__51.MoveNext() in D:\repo\projects\xxx\mainapi\xxx.Main.Api\Messaging\MessageBus.cs:line 34 

StartSendActivity 看起来像简单的诊断方法,我无法确定,那里可能会失败。

【问题讨论】:

    标签: asp.net-core amazon-sqs masstransit


    【解决方案1】:

    您不能通过 MassTransit 发送 object 或其他类似的类型。消息必须是引用类型。

    Relevant Documentation

    【讨论】:

    • 所以我不能使用dynamic 作为消息?那么为什么它适用于Bus.Factory.CreateUsingAmazonSqs(见上文)?我可以使用什么方法为所有消息类型设置CorrelationId
    • 不知道,你的堆栈跟踪太短了。而且我不支持建议之外的消息类型 -​​ 所以我建议按照提供的链接中的指导进行操作。
    • 我添加了完整的堆栈跟踪。无论如何,如果我按照文档的建议切换到匿名类型 - 是否有任何方法可以为所有消息类型设置 CorrelationId
    • 是的,但我应该为每个Publish 调用都这样做,我试图避免这种情况。不管怎样,谢谢你们的cmets。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-23
    • 1970-01-01
    • 1970-01-01
    • 2018-10-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多