消息复制任务和应用程序

如消息复制和跨区域联合一文中所述,在服务总线实体对之间以及服务总线与其他消息源和目标之间的消息序列复制通常依赖于 Azure Functions。

Azure Functions 是一种可缩放且可靠的执行环境,用于配置和运行无服务器应用程序,包括消息复制和联合任务。

在本概述中,你将了解 Azure Functions 针对此类应用程序提供的内置功能、可为转换任务调整和修改的代码块,以及如何配置 Azure Functions 应用程序,使其与服务总线和其他 Azure 消息传送服务完美集成。 有关更多详细信息,本文将指向 Azure Functions 文档。

什么是复制任务?

复制任务是从源接收事件并将其转发到目标。 大多数复制任务转发事件不变。 如果源和目标使用不同的协议,任务还可以映射其元数据结构。

复制任务通常是无状态的。 它们不会在顺序或并行执行中共享状态或其他副作用。 批处理和链式处理也可以利用流的现有状态。

这一特性使复制任务区别于聚合任务,后者通常是有状态的,并由分析框架和服务(如 Azure 流分析)支持。

Azure Functions 中的复制应用程序和任务

在 Azure Functions 中,复制任务使用触发器接收源端的消息。 该任务使用 输出绑定 或相应的 Azure 客户端库将副本发送到目标。

触发器 输出
Azure 事件中心触发器 Azure 事件中心输出绑定
Azure 服务总线触发器 Azure 服务总线输出绑定
Azure IoT 中心触发器 Azure IoT 中心输出绑定
Azure 事件网格触发器 Azure 事件网格输出绑定
Azure 队列存储触发器 Azure 队列存储输出绑定
Apache Kafka 触发器 Apache Kafka 输出绑定
RabbitMQ 触发器 RabbitMQ 输出绑定
Azure 通知中心输出绑定
Azure SignalR 服务 output binding
Twilio SendGrid 输出绑定

在新的.NET复制应用中使用Azure Functions .NET孤立工作者模型。 .NET 在进程模型的支持将于2026年11月10日结束。 要更新现有的进程内复制应用,请按照 .NET 隔离工作进程迁移指南。

你可以将多个复制任务部署到同一个函数应用。 通过 Azure Functions Premium,多个功能应用可以共享同一个应用服务套餐。 当集成需要特定语言库时,这种配置还允许你将用不同语言编写的复制任务共置。

如果有面向批处理的触发器可用,请优先使用。 接收完整的事件或消息结构,而不是依赖 Azure Functions 的绑定表达式,从而使任务能够保留源元数据。

每个函数应以它连接的源和目标命名。 连接和命名空间设置时,使用同一个名称作为前缀。

数据和元数据映射

映射目标支持的消息正体和应用属性。 目标会为经纪人拥有的字段(如排队时间和序列号)分配新值。 在 repl-sequence 和 repl-enqueue-time 应用程序属性中保留源值。 如果已有任一属性,则用分号分隔符添加新值。 不要复制送货次数或锁定信息。

独立辅助角色工作进程的 服务总线 输出绑定支持简单输出类型,但不支持 ServiceBusMessage。 当复制任务必须保留 服务总线 消息元数据时,请使用 ServiceBusClient。 使用客户端还能让你的代码处理并记录发送失败。 以下示例使用两个目标的Microsoft客户端库来明确此行为。

重试策略

根据源触发器和目标客户端来配置重试设置。 对于Event Hubs的触发器,你可以应用Azure Functions函数级的重试策略。 Event Hubs 在该执行的重试策略完成之前,不会写入检查点。

服务总线 触发器使用实体和 host.json中配置的重试和死符行为。 不要给 服务总线 触发器应用函数级重试属性。 失败批次会被完整重新送达,这可能导致目标处产生重复。 当消息超过源实体的 MaxDeliveryCount时,服务总线 会将其移至死符队列。 将 MaxDeliveryCount 设置为所需的恢复时间窗口,并监控死信队列。

Microsoft客户端库还应用其配置的重试策略来发送操作。 如果发送仍然失败,让异常从函数中逃逸出来,这样源触发器可以重新尝试批次。 将重试次数设为无限的事件中心重试策略会暂停受影响分区的检查点进度,直到目标发送操作成功。

建立一个复制应用主机

复制应用是一种 Azure Functions 应用,用于托管一个或多个复制任务。 使用支持的 Functions 4.x 运行时版本和 .NET 隔离工作者模型。 根据你的规模和网络需求,创建应用。

为函数应用使用 系统分配的或用户分配的托管身份。 授予身份从每个源接收并发送给每个目标的权限。 为触发器配置基于身份的连接,并为Azure客户端注册提供完全限定的目标命名空间。 对于示例,设 telemetrySourceConnection__fullyQualifiedNamespace、 telemetryTarget__fullyQualifiedNamespace、 jobsTransferSourceConnection__fullyQualifiedNamespace和 jobsTransferTarget__fullyQualifiedNamespace。

以下 Program.cs 示例通过使用 DefaultAzureCredential来注册目标服务连接。 在 Azure 中,DefaultAzureCredential 使用函数应用的托管标识。

using Azure.Identity;
using Microsoft.Azure.Functions.Worker.Builder;
using Microsoft.Extensions.Azure;
using Microsoft.Extensions.Hosting;

var builder = FunctionsApplication.CreateBuilder(args);

builder.Services.AddAzureClients(clientBuilder =>
{
    clientBuilder.AddEventHubProducerClientWithNamespace(
        builder.Configuration["telemetryTarget:fullyQualifiedNamespace"],
        "telemetry-copy");
    clientBuilder.AddServiceBusClientWithNamespace(
        builder.Configuration["jobsTransferTarget:fullyQualifiedNamespace"]);
    clientBuilder.UseCredential(new DefaultAzureCredential());
});

builder.Build().Run();

项目需要隔离工作进程包、Event Hubs 和 服务总线 绑定扩展、Azure.Identity、Azure.Messaging.EventHubs、Azure.Messaging.ServiceBus 和 Microsoft.Extensions.Azure。 有关当前包和目标框架要求,请参见 .NET 隔离工作进程指南 和 迁移指南。

通过虚拟网络访问Event Hubs或服务总线命名空间的复制应用必须使用支持虚拟网络集成的托管计划。 有关详细信息,请参阅 Azure Functions 网络选项。

示例

以下 .NET 独立辅助角色示例会成批接收消息,复制由客户控制的消息正文和元数据,并使用已在 Program.cs 中注册的目标客户端。

要在事件中心之间复制事件数据,请使用事件中心触发器和一个 EventHubProducerClient:

using System;
using System.Collections.Generic;
using System.Globalization;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Producer;
using Microsoft.Azure.Functions.Worker;

public sealed class EventHubReplicationFunction
{
    private readonly EventHubProducerClient outputClient;

    public EventHubReplicationFunction(EventHubProducerClient outputClient)
    {
        this.outputClient = outputClient;
    }

    [Function("telemetry")]
    [ExponentialBackoffRetry(-1, "00:00:05", "00:15:00")]
    public async Task Telemetry(
        [EventHubTrigger(
            "telemetry",
            ConsumerGroup = "%telemetrySourceConsumerGroup%",
            Connection = "telemetrySourceConnection")]
        EventData[] input,
        CancellationToken cancellationToken)
    {
        foreach (IGrouping<string, EventData> group in input.GroupBy(item => item.PartitionKey))
        {
            var options = new CreateBatchOptions { PartitionKey = group.Key };
            EventDataBatch batch = await outputClient.CreateBatchAsync(options, cancellationToken);

            try
            {
                foreach (EventData item in group)
                {
                    EventData copy = CopyEvent(item);

                    if (!batch.TryAdd(copy))
                    {
                        if (batch.Count == 0)
                        {
                            throw new InvalidOperationException("An event exceeds the maximum Event Hubs batch size.");
                        }

                        await outputClient.SendAsync(batch, cancellationToken);
                        batch.Dispose();
                        batch = await outputClient.CreateBatchAsync(options, cancellationToken);

                        if (!batch.TryAdd(copy))
                        {
                            throw new InvalidOperationException("An event exceeds the maximum Event Hubs batch size.");
                        }
                    }
                }

                if (batch.Count > 0)
                {
                    await outputClient.SendAsync(batch, cancellationToken);
                }
            }
            finally
            {
                batch.Dispose();
            }
        }
    }

    private static EventData CopyEvent(EventData input)
    {
        var output = new EventData(input.EventBody)
        {
            ContentType = input.ContentType,
            CorrelationId = input.CorrelationId,
            MessageId = input.MessageId,
        };

        foreach (KeyValuePair<string, object> property in input.Properties)
        {
            output.Properties[property.Key] = property.Value;
        }

        AppendReplicationProperty(
            output.Properties,
            "repl-enqueue-time",
            FormatSystemProperty(input, "x-opt-enqueued-time"));
        AppendReplicationProperty(
            output.Properties,
            "repl-sequence",
            FormatSystemProperty(input, "x-opt-sequence-number"));

        return output;
    }

    private static string FormatSystemProperty(EventData input, string name)
    {
        if (!input.SystemProperties.TryGetValue(name, out object? value))
        {
            throw new InvalidOperationException($"Event Hubs didn't provide {name}.");
        }

        return value switch
        {
            DateTime dateTime => dateTime.ToUniversalTime().ToString("O"),
            DateTimeOffset dateTimeOffset => dateTimeOffset.ToUniversalTime().ToString("O"),
            IFormattable formattable => formattable.ToString(null, CultureInfo.InvariantCulture),
            _ => value.ToString()!,
        };
    }

    private static void AppendReplicationProperty(
        IDictionary<string, object> properties,
        string name,
        string value)
    {
        properties[name] = properties.TryGetValue(name, out object? existing)
            ? $"{existing};{value}"
            : value;
    }
}

若要在 服务总线 实体之间复制消息,请使用 服务总线 触发器和一个 ServiceBusClient:

using System;
using System.Collections.Generic;
using System.Globalization;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Azure.Messaging.ServiceBus;
using Microsoft.Azure.Functions.Worker;

public sealed class ServiceBusReplicationFunction
{
    private readonly ServiceBusSender sender;

    public ServiceBusReplicationFunction(ServiceBusClient client)
    {
        sender = client.CreateSender("jobs");
    }

    [Function("jobs-transfer")]
    public async Task JobsTransfer(
        [ServiceBusTrigger(
            "jobs-transfer",
            Connection = "jobsTransferSourceConnection",
            IsBatched = true,
            IsSessionsEnabled = true)]
        ServiceBusReceivedMessage[] input,
        CancellationToken cancellationToken)
    {
        foreach (IGrouping<(string SessionId, string PartitionKey), ServiceBusReceivedMessage> group
            in input.GroupBy(item => (item.SessionId, item.PartitionKey)))
        {
            ServiceBusMessageBatch batch = await sender.CreateMessageBatchAsync(cancellationToken);

            try
            {
                foreach (ServiceBusReceivedMessage item in group)
                {
                    ServiceBusMessage copy = CopyMessage(item);

                    if (!batch.TryAddMessage(copy))
                    {
                        if (batch.Count == 0)
                        {
                            throw new InvalidOperationException("A message exceeds the maximum Service Bus batch size.");
                        }

                        await sender.SendMessagesAsync(batch, cancellationToken);
                        batch.Dispose();
                        batch = await sender.CreateMessageBatchAsync(cancellationToken);

                        if (!batch.TryAddMessage(copy))
                        {
                            throw new InvalidOperationException("A message exceeds the maximum Service Bus batch size.");
                        }
                    }
                }

                if (batch.Count > 0)
                {
                    await sender.SendMessagesAsync(batch, cancellationToken);
                }
            }
            finally
            {
                batch.Dispose();
            }
        }
    }

    private static ServiceBusMessage CopyMessage(ServiceBusReceivedMessage input)
    {
        var output = new ServiceBusMessage(input.Body)
        {
            ContentType = input.ContentType,
            CorrelationId = input.CorrelationId,
            MessageId = input.MessageId,
            PartitionKey = input.PartitionKey,
            ReplyTo = input.ReplyTo,
            ReplyToSessionId = input.ReplyToSessionId,
            SessionId = input.SessionId,
            Subject = input.Subject,
            To = input.To,
        };

        foreach (KeyValuePair<string, object> property in input.ApplicationProperties)
        {
            output.ApplicationProperties[property.Key] = property.Value;
        }

        AppendReplicationProperty(
            output.ApplicationProperties,
            "repl-enqueue-time",
            input.EnqueuedTime.ToString("O"));
        AppendReplicationProperty(
            output.ApplicationProperties,
            "repl-sequence",
            input.SequenceNumber.ToString(CultureInfo.InvariantCulture));

        return output;
    }

    private static void AppendReplicationProperty(
        IDictionary<string, object> properties,
        string name,
        string value)
    {
        properties[name] = properties.TryGetValue(name, out object? existing)
            ? $"{existing};{value}"
            : value;
    }
}

服务总线 示例通过会话 ID 和分区键分组消息,然后按源顺序发送每个组。 这种方法保留了会话之间的相对顺序。 这两个例子都创建了规模受限的批次,且这些批次都保持在目标服务的最大批处理大小范围内。 服务总线 的例子不会复制源的存活时间,因为复制会重启目标的生命周期。 对于过期感知复制,计算剩余寿命或应用目标实体的过期策略。

对于生产复制路径,还要决定如何处理重复传递、消息过期、事务、调度消息、死符消息以及目标协议无法表示的元数据。

监控

使用Azure Functions Monitoring来监控复制应用。

Application Insights 应用映射 可视化了复制任务与其源和目标之间的依赖关系。 实时指标 在函数应用运行时提供低延迟的诊断信息。

后续步骤