【问题标题】:Azure WebJobs SDK: Not calling message.complete() on a specific ServiceBus triggered functionAzure WebJobs SDK:不在特定的 ServiceBus 触发函数上调用 message.complete()
【发布时间】:2016-09-27 17:30:35
【问题描述】:

我有一个如下形式的服务总线触发函数:

public static void ProcessJobs([ServiceBusTrigger("topicname", "subscriptionname", AccessRights.Listen)] BrokeredMessage input, [ServiceBus("topic2name", "subscription2name", AccessRights.Send)] out BrokeredMessage output)
{
    output = new BrokeredMessage(input.GetBody<String>());
}

我的用例是这个函数只是从一个主题中获取消息并将它们推送到另一个主题。我不想在进程中从源主题中删除消息。

这可以实现吗?

此外,我在哪里可以找到有关 AccessRights 以及它们如何影响消息访问的更多信息。例如:在上面的示例中,我使用 AccessRights.Listen 从输入主题获取消息,但是一旦对这些消息执行函数调用,它似乎仍然是“删除”消息。

【问题讨论】:

  • 你为什么不想完成消息?
  • @Thomas:不确定这是否是反模式,但我正在考虑嵌套服务总线队列。所以在这种情况下,我有一个父队列将其项目移交给子队列。一旦孩子完成,父队列项目(连同子队列项目)被移除。这更多是出于我的应用程序的语义而不是出于特定的技术原因。
  • 所以如果你使用一个主题,你可以有多个订阅并使用过滤器来处理父/子队列之间的通信。如果您未完成该消息,它也将可用,因此您将再次对其进行重新处理。

标签: azure azure-webjobssdk azure-servicebus-topics


【解决方案1】:

正如 Jambor 所说,默认行为是在函数成功完成时完成消息(参见documentation),如果函数失败则放弃。

你可以在SDK Repo查看MessageProcessor类的代码看到这个行为的实现:

public virtual async Task CompleteProcessingMessageAsync(BrokeredMessage message, FunctionResult result, CancellationToken cancellationToken)
{
    if (result.Succeeded)
    {
        if (!MessageOptions.AutoComplete)
        {
            // AutoComplete is true by default, but if set to false
            // we need to complete the message
            cancellationToken.ThrowIfCancellationRequested();
            await message.CompleteAsync();
        }
    }
    else
    {
        cancellationToken.ThrowIfCancellationRequested();
        await message.AbandonAsync();
    }
}

有趣的一点:这是一个虚函数。

ServiceBusConfiguration 公开了一个 MessagingProvider 属性。

如果您查看SDK Repo 中默认MessagingProvider 类的代码,您会发现您可以重写负责创建新MessageProcessor 的方法:

/// <summary>
/// Creates a <see cref="MessageProcessor"/> for the specified ServiceBus entity.
/// </summary>
/// <param name="entityPath">The ServiceBus entity to create a <see cref="MessageProcessor"/> for.</param>
/// <returns>The <see cref="MessageProcessor"/>.</returns>
public virtual MessageProcessor CreateMessageProcessor(string entityPath)
{
    if (string.IsNullOrEmpty(entityPath))
    {
        throw new ArgumentNullException("entityPath");
    }
    return new MessageProcessor(_config.MessageOptions);
}

这个函数也是虚拟的。

现在您可以创建自己的 MessagingProviderMessageProcessor 实现:

public class CustomMessagingProvider : MessagingProvider
{
    private readonly ServiceBusConfiguration _config;

    public CustomMessagingProvider(ServiceBusConfiguration config) : base(config)
    {
        _config = config;
    }

    public override MessageProcessor CreateMessageProcessor(string entityPath)
    {
        if (string.IsNullOrEmpty(entityPath))
        {
            throw new ArgumentNullException("entityPath");
        }
        return new CustomMessageProcessor(_config.MessageOptions);
    }

    class CustomMessageProcessor : MessageProcessor
    {
        public CustomMessageProcessor(OnMessageOptions messageOptions) : base(messageOptions)
        {
        }

        public override async Task CompleteProcessingMessageAsync(BrokeredMessage message, FunctionResult result, CancellationToken cancellationToken)
        {
            if (!result.Succeeded)
            {
                cancellationToken.ThrowIfCancellationRequested();
                await message.AbandonAsync();
            }
        }
    }
}

然后像这样配置你的 JobHost:

public static void Main()
{
    var config = new JobHostConfiguration();
    var sbConfig = new ServiceBusConfiguration
    {
        MessageOptions = new OnMessageOptions
        {
            AutoComplete = false
        }
    };
    sbConfig.MessagingProvider = new CustomMessagingProvider(sbConfig);
    config.UseServiceBus(sbConfig);
    var host = new JobHost(config);
    host.RunAndBlock();
}

这是针对此问题的技术部分...

现在,如果您没有完成您的消息,该消息将一次又一次地用于相同的功能,直到您到达MaxDeliveryCount,然后您的消息将是死信。因此,即使将您的函数设计为幂等的,我也很确定这不是您想要的。

也许你应该多解释一下你想要达到的目标?

如果您正在寻找父子队列通信(请参阅问题 cmets),有一篇很好的文章解释了如何使用 ASB 设计工作流:

否则,你可以看看 BrokerMessage 对象上的Defer 方法:

表示接收者希望推迟处理此消息。

它将允许您在处理子消息之前保留父消息。

【讨论】:

    【解决方案2】:

    我不想在进程中从源主题中删除消息。

    this document我们可以知道,ServiceBusTrigger会在触发完成后调用完成函数。这是默认模式。这是那篇文章中的 sn-p:

    SDK 在 PeekLock 模式下收到一条消息,如果函数成功完成,则对该消息调用 Complete,如果函数失败,则调用 Abandon。如果函数运行时间超过 PeekLock 超时时间,锁会自动更新。

    【讨论】:

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