【问题标题】:RabbitMq - ConversationId vs CorrelationId - Which is the more appropriate for tracking a specific request?RabbitMq - ConversationId 与 CorrelationId - 哪个更适合跟踪特定请求?
【发布时间】:2019-03-21 16:42:37
【问题描述】:

RabbitMQ 似乎有两个非常相似的属性,我并不完全理解其中的区别。 ConversationIdCorrelationId

我的用例如下。我有一个生成Guid 的网站。该网站调用一个 API,将该唯一标识符添加到 HttpRequest 标头中。这反过来将消息发布到 RabbitMQ。该消息由第一个消费者处理并在其他地方传递给另一个消费者,依此类推。

出于记录目的,我想记录一个将初始请求与所有后续操作联系在一起的标识符。这对于整个应用程序不同部分的旅程来说应该是独一无二的。因此。当登录到像 Serilog/ElasticSearch 这样的东西时,就可以很容易地看到哪个请求触发了初始请求,并且整个应用程序中该请求的所有日志条目都可以关联在一起。

我创建了一个提供程序,它查看传入的HttpRequest 以获取标识符。我将其称为“CorrelationId”,但我开始怀疑是否真的应该将其命名为“ConversationId”。就 RabbitMQ 而言,“ConversationId”的概念更适合这个模型,还是“CorrelationId”更好?

这两个概念有什么区别?

在代码方面,我希望执行以下操作。首先在我的 API 中注册总线并配置 SendPublish 以使用来自提供商的 CorrelationId

// bus registration in the API
var busSettings = context.Resolve<BusSettings>();
// using AspNetCoreCorrelationIdProvider
var correlationIdProvider = context.Resolve<ICorrelationIdProvider>();

var busControl = Bus.Factory.CreateUsingRabbitMq(cfg =>
{
    cfg.Host(
        new Uri(busSettings.HostAddress),
        h =>
        {
            h.Username(busSettings.Username);
            h.Password(busSettings.Password);
        });
    cfg.ConfigurePublish(x => x.UseSendExecute(sendContext =>
    {
        // which one is more appropriate
        //sendContext.ConversationId = correlationIdProvider.GetCorrelationId();
        sendContext.CorrelationId = correlationIdProvider.GetCorrelationId();
    }));
});

作为参考,这是我的简单提供者界面

// define the interface
public interface ICorrelationIdProvider
{
    Guid GetCorrelationId();
}

还有 AspNetCore 实现,它提取调用客户端(即网站)设置的唯一 ID。

public class AspNetCoreCorrelationIdProvider : ICorrelationIdProvider
{
    private IHttpContextAccessor _httpContextAccessor;

    public AspNetCoreCorrelationIdProvider(IHttpContextAccessor httpContextAccessor)
    {
        _httpContextAccessor = httpContextAccessor;
    }

    public Guid GetCorrelationId()
    {
        if (_httpContextAccessor.HttpContext.Request.Headers.TryGetValue("correlation-Id", out StringValues headers))
        {
            var header = headers.FirstOrDefault();
            if (Guid.TryParse(header, out Guid headerCorrelationId))
            {
                return headerCorrelationId;
            }
        }

        return Guid.NewGuid();
    }
}

最后,我的服务主机是简单的 Windows 服务应用程序,它们坐下来使用已发布的消息。他们使用以下内容来获取 CorrelationId,并且很可能会在其他服务主机中发布给其他消费者。

public class MessageContextCorrelationIdProvider : ICorrelationIdProvider
{
    /// <summary>
    /// The consume context
    /// </summary>
    private readonly ConsumeContext _consumeContext;

    /// <summary>
    /// Initializes a new instance of the <see cref="MessageContextCorrelationIdProvider"/> class.
    /// </summary>
    /// <param name="consumeContext">The consume context.</param>
    public MessageContextCorrelationIdProvider(ConsumeContext consumeContext)
    {
        _consumeContext = consumeContext;
    }

    /// <summary>
    /// Gets the correlation identifier.
    /// </summary>
    /// <returns></returns>
    public Guid GetCorrelationId()
    {
        // correlationid or conversationIs?
        if (_consumeContext.CorrelationId.HasValue && _consumeContext.CorrelationId != Guid.Empty)
        {
            return _consumeContext.CorrelationId.Value;
        }

        return Guid.NewGuid();
    }
}

然后我的消费者中有一个记录器,它使用该提供程序来提取CorrelationId

public async Task Consume(ConsumeContext<IMyEvent> context)
{
    var correlationId = _correlationProvider.GetCorrelationId();
    _logger.Info(correlationId, $"#### IMyEvent received for customer:{context.Message.CustomerId}");

    try
    {
        await _mediator.Send(new SomeOtherRequest(correlationId) { SomeObject: context.Message.SomeObject });
    }
    catch (Exception e)
    {
        _logger.Exception(e, correlationId, $"Exception:{e}");
        throw;
    }

    _logger.Info(correlationId, $"Finished processing: {DateTime.Now}");
}

阅读docs,它说以下关于“ConversationId”:

对话是由发送的第一条消息创建的,或者 已发布,其中没有可用的现有上下文(例如,当 消息通过使用 IBus.Send 或 IBus.Publish 发送或发布)。如果 现有上下文用于发送或发布消息, ConversationId 被复制到新消息中,确保一组 同一对话中的消息具有相同的标识符。

现在我开始认为我的术语混淆了,从技术上讲,这是一次对话(尽管“对话”就像“电话游戏”)。

那么,在这个用例中是CorrelationId,还是ConversationId?请帮我弄好我的术语!

【问题讨论】:

  • 似乎 ConversationId 是 MassTransit/NServiceBus 的东西?我在 RabbitMQ 或 AMQP 文档中没有找到任何提及。 CorrelationId 一个特定的 AMQP 事物。规范说“没有正式行为,但在请求消息中使用时可能包含私有响应队列的名称”,并且在执行 RPC 请求时通常用于响应队列名称。所以我认为您的用例是 ConversationId,而不是 CorrelationId。
  • ConversationId 仅限于单个消息序列,您在其中收到一条消息并通过多个消费者。 CorrelationId 更长寿,一个关联 id 可以进行多次对话。
  • 除此之外,对话 ID 是自动生成的,相关 ID 是任意的,您可以从您的域中获取。例如,在我们的 ShoppingCart 传奇中,我们使用订单 ID 作为关联 ID。
  • 我在 ConversationId 上得到纠正。当然,如果您想保留 CorrelationId 用于执行 RPC 的正常目的,您可以添加自己的标头。
  • 当我遇到这个问题时,我决定使用我的业务领域中的一个术语。它最终只是被称为 lineId 而我没有使用标题。我只是使用了我自己的自定义correlationId,这很有帮助,因为不同的rabbitmq客户端可以决定对这些标头做不同的事情。如果我想不那么健谈,我还可以在每条消息中包含多个 lineid。

标签: c# rabbitmq masstransit


【解决方案1】:

在消息对话中(暗示乐谱),可以有一条消息(我告诉你做某事,或者我告诉正在听的每个人发生了某事)或多条消息(我告诉你做某事,然后你告诉了其他人,或者我告诉所有在听的人发生了什么事,而这些听众告诉了他们的朋友,等等)。

使用 MassTransit,从第一条消息到最后一条消息,如果使用得当,这些消息中的每一条都将具有相同的 ConversationId。 MassTransit 在消息消费期间将属性从ConsumeContext 复制到每条传出消息,未经修改。这使得所有内容都成为同一 trace 的一部分 - 一次对话。

但是,MassTransit 默认不设置 CorrelationId。如果消息属性名为 CorrelationId(或 CommandId,或 EventId),则可以自动设置它,或者您也可以添加自己的名称。

如果 CorrelationId 出现在消费消息中,则任何传出消息都会将该 CorrelationId 属性复制到 InitiatorId 属性(因果关系 - 消费消息启动了后续消息的创建)。这形成了一个链(或跟踪术语中的跨度),可以跟踪该链以显示来自初始消息的消息的传播。

应该将 CorrelationId 视为命令或事件的标识符,以便可以在整个系统日志中看到该命令的效果。

在我看来,您来自 HTTP 的输入可能是 Initiator,因此将该标识符复制到 InitiatorId 并为消息创建一个新的 CorrelationId,或者您可能只想对初始 CorrelationId 使用相同的标识符并让后续消息使用它作为发起者。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多