Microsoft Orleans Durable Jobs 完全指南:分布式一次性任务的调度、持久化与迁移
后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载本指南以src/Orleans.DurableJobs/README.md为骨架系统讲解 Microsoft Orleans 的 Durable Jobs持久化任务模块如何安装与配置内存版/Azure Blob 版存储、如何通过ILocalDurableJobManager调度与取消一次性任务、如何深入理解基于时间分片Shard的分布式执行与故障转移原理以及如何借助命名 Journal Provider 平滑完成存储迁移与退役。读完本篇你将能够在自己的 Orleans Silo 中落地预约提醒、延迟处理、定时工作流步骤等真实场景并掌握配置调优、优雅关闭与 A→B 迁移的完整实操方法。Durable Jobs 是什么与 Reminders 的定位差异Microsoft Orleans Durable Jobs 提供了一套分布式、可扩展的一次性任务调度系统让任务在指定时间点被精确执行。与面向周期性重复任务的 Orleans Reminders 不同Durable Jobs 专为未来某个时刻只发生一次的事件设计典型场景包括预约通知appointment notifications延迟处理delayed processing定时工作流步骤scheduled workflow steps时间触发的触发器time-based triggers其核心能力来自 README.md包括至少执行一次任务保证至少被调度执行一次At Least Once持久化任务在 Grain 失活、Silo 重启后依然存活分布式任务被自动分发并在各 Silo 之间重新平衡可靠失败任务可按可配置策略自动重试丰富元数据每个任务可携带自定义元数据可持久化取消取消请求会阻止后续尝试启动但已在运行中的一次尝试仍可能完整执行完。从源码看任务的调度由LocalDurableJobManagerLocalDurableJobManager.cs作为 System Target 提供它同时实现了ILocalDurableJobManager、ILocalDurableJobManagerSystemTarget与ILifecycleParticipantISiloLifecycle也就是说它既是本地调度入口也能通过 System Target 消息把取消请求路由到持有任务分片的远端 Silo还是挂载在 Silo 生命周期Active 阶段启动、停止阶段优雅关闭上的参与者。快速开始安装与基础配置安装 NuGet 包dotnet add package Microsoft.Orleans.DurableJobs生产环境需要持久化时再安装对应的存储 Provider以 Azure Storage 为例dotnet add package Microsoft.Orleans.DurableJobs.AzureStorage内存存储配置开发/测试内存版依赖AddVolatileJournalStorage()组合UseJournaledDurableJobs()源码见 DurableJobsExtensions.cs 中的UseInMemoryDurableJobsusing Microsoft.Extensions.Hosting; using Orleans.Hosting; var builder Host.CreateApplicationBuilder(args); builder.UseOrleans(siloBuilder { siloBuilder .UseLocalhostClustering() // Configure in-memory Durable Jobs (no persistence) .UseInMemoryDurableJobs(); }); await builder.Build().RunAsync();注意易失性存储Volatile只把任务保留在进程生命周期内Silo 重启即丢失仅适用于开发与测试。Azure Blob 存储配置生产using Microsoft.Extensions.Hosting; using Orleans.Hosting; var builder Host.CreateApplicationBuilder(args); builder.UseOrleans(siloBuilder { siloBuilder .UseLocalhostClustering() // Configure Azure Storage Durable Jobs .UseAzureBlobDurableJobs(options { options.BlobServiceClient new Azure.Storage.Blobs.BlobServiceClient(YOUR_CONNECTION_STRING); options.ContainerName durable-jobs; }); }); await builder.Build().RunAsync();Azure 版支持UseAzureBlobDurableJobsappend-blob WAL block-blob checkpoint与UseAzureTableDurableJobsTable journal headers data generations两种后端详细配置连接字符串、托管身份DefaultAzureCredential、按环境区分容器名等可参考 Azure Storage 版 README。生产环境推荐使用 Managed Identity 而非连接字符串。高级配置builder.UseOrleans(siloBuilder { siloBuilder .UseLocalhostClustering() .UseInMemoryDurableJobs() .ConfigureServices(services { services.ConfigureDurableJobsOptions(options { // Duration of each job shard (jobs are partitioned by time) options.ShardDuration TimeSpan.FromMinutes(5); // Load eligible shards within this horizon and check at this interval options.ShardLoadLookaheadPeriod TimeSpan.FromMinutes(10); options.ShardCheckInterval TimeSpan.FromMinutes(5); // Maximum number of jobs that can execute concurrently on each silo options.MaxConcurrentJobsPerSilo 100; // Custom retry policy options.ShouldRetry (context, exception) { // Retry up to 3 times with exponential backoff if (context.DequeueCount 3) { var delay TimeSpan.FromSeconds(Math.Pow(2, context.DequeueCount)); return DateTimeOffset.UtcNow.Add(delay); } return null; // Dont retry }; }); }); });选择命名 Journal 存储多 Provider 并存在 README.md 的 Select named journal storage 一节中核心结论是UseJournaledDurableJobs同时提供ISiloBuilder与IServiceCollection两个重载独立于存储后端安装 journaled shard manager 与 Durable Jobs JSON 元数据支持。你需要先注册一个支持 Catalog 的 Journal Provider再用DurableJobsOptions.ActiveProviderName选中它默认名称为DefaultDrainingProviderNames初始为空。siloBuilder .AddAzureBlobJournalStorage(jobs-a, options { options.BlobServiceClient originalClient; options.ContainerName jobs-original; }) .AddAzureBlobJournalStorage(jobs-b, options { options.BlobServiceClient currentClient; options.ContainerName jobs-current; }) .UseJournaledDurableJobs(options { options.ActiveProviderName jobs-b; options.DrainingProviderNames.Add(jobs-a); });关键约束可在源码 DurableJobsOptions.cs 中印证命名的 Blob、Table、Redis、S3 与 Volatile 注册会把存储、Catalog、State-Manager 工厂绑定到同一个命名空间。显式命名后默认 Provider 仍可独立用于 Grain journaling。每个物理命名空间只能用一个名字注册UseInMemoryDurableJobs、UseAzureBlobDurableJobs、UseAzureTableDurableJobs只是默认绑定上的便捷组合。Default绑定沿用既有的未命名后端选项管道非默认绑定如jobs-a使用命名后端选项。大多数部署只选一个 Provider无缓存的已知 Shard 查找会直接读取该 Provider 的元数据与其他为 Grain 状态注册的 Journal Provider 无关。迁移期间发现discovery覆盖写 Provider 与 draining Provider无缓存的已知 ID 查找先读写 Provider必要时再查 draining Provider。Provider 不可用会表现为查找失败缺失则要求对所有选中 Provider 都检查成功。新调度使用在写 Provider中创建的分片已存在的分片在其整个生命周期内保留原 Provider 负责所有权、重放、执行、重试、重调度、取消、压缩与删除。Timestamp/GUID 形式的分片 ID、jobs/shards/id路径以及 Provider 无关的序列化DurableJob句柄均不改变——一个分片在整个生命周期只有一个权威存储位置。切换并退役一个 ProviderA→B 迁移绑定在启动时固定。按 README 给出的步骤执行 A→B 迁移预先部署在任何一个 Silo 创建 B 的工作之前在每个调度 Silo 上同时部署两个绑定并让 A 负责写、B 负责 draining。滚动切换随后将 B 滚动为写 Provider、A 变为 draining。此滚动期间A 配置的 Silo 仍可创建 A 分片。当每个调度 Silo 都使用 B且之前的调度调用都已结束时切换完成。若需要严格切换边界可在变更期间暂停应用调度。保留 A 的变更权限draining 需要列出listing、读取、认领claims、持久化更新、重试、取消、压缩与删除。清点检查所有写方切换完成后用IDurableJobsStorageInspector.InspectAsync(jobs-a, cancellationToken)检查 A 的完整命名空间解决所有剩余分片包括未来日期、已拥有、中毒、无法识别的条目并按后端实时列举行为重复成功盘点。确认退役只有经过验证的完全 drain 之后才把 A 从DrainingProviderNames移除只有当其他消费者也完成后才退役该绑定/存储。检查器返回的DurableJobsStorageStatus字段定义见 IDurableJobsStorageInspector.cs包含ProviderNameJournal Provider 名称IsWriteProvider该 Silo 是否用它创建新分片ShardCount分片命名空间中的 Journal 数量OwnedShardCount有记录所有者的已识别分片数PoisonedShardCount中毒分片数UnrecognizedShardCount元数据无法解析的 Journal 数OldestShardStartTime/NewestShardStartTime最早/最晚已识别分片开始时间可为 null。需要特别说明检查失败会令调用出错此时应显示unknown计数描述的是分片每个分片可包含大量任务owned/poisoned 计数可能重叠退役要求成功完全清零清单 集群范围内写方全部切换完成。回滚保留两个 Provider 均被选中把写切回 A、B 继续 drainingB 中已有任务继续留在 B。整个迁移期间应保留存储命名空间与兼容的 Journal 读取器。仓库中还提供了两个配套资料重启式迁移示例包含基于 Azurite 的持久化、保留句柄的取消与本地清单检查以及迁移部署与监控指南。Shard 发现与前瞻Discovery Lookahead每个 Silo 会发现在其 Durable Jobs 时间 Provider 时钟的ShardLoadLookaheadPeriod范围内开始时间的分片默认前瞻 10 分钟。一个被发现的 Shard 在开始时间进入ShardActivationBufferPeriod后才会启动处理默认 5 分钟缓冲避免长时间持有空闲分片。ShardCheckInterval控制周期性发现与可写分片清理默认 5 分钟集群成员变化也会触发检查。约束细节源码 DurableJobsOptions.cs 与校验器可印证lookahead 接受非负时长零表示选择开始时间不晚于当前时间的分片发现视界上限为DateTimeOffset.MaxValuecheck interval 合法区间为 1 到 4294967294 毫秒即uint.MaxValue - 1。分片 Journal 的命名形如jobs/shards/20260909T1200000000000Z-unique-id固定宽度的 UTC 开始时间放在唯一后缀之前因此序数名称顺序就是分片开始时间顺序。每次扫描用包含JournalCatalogListOptions.MaxId覆盖前瞻视界的列表列出jobs/shards/前缀范围包含所有更早的开始时间——包括长时间宕机后逾期积压的任务。分片认领机制见 LocalDurableJobManager.cs 中ProcessShardDiscoveryAsync/ComputeClaimBudget每次周期或成员检查都启动一次全新的、本地作用域的扫描覆盖所有选中的 Provider候选者共享最旧优先oldest-first排序与聚合认领预算claim budget认领按最旧优先执行且要求使用快照的 ETag 做条件更新因此并发的所有权变化会让过期认领被拒绝Provider 必须提供元数据 ETag 并强制条件更新缺失 ETag 会表现为发现错误分配到的分片在被打开时即交付后续候选仍在评估时执行就可以开始预算用尽后新认领停止但本地已拥有的分片仍保持合格。同时有孤儿分片认领慢启动ramp-up机制防止灾后恢复时新启动的 Silo 一次性认领过多分片ShardClaimInitialBudget默认 2在启动时限制立即认领的孤儿分片数随运行时间线性增长到ShardClaimMaxBudget默认 20经过ShardClaimRampUpDuration默认 5 分钟后限制完全移除Silo 过载时新认领暂停预算返回 0。关于不同后端Catalog Provider 利用各自的存储能力施加原始前缀与范围约束Azure Table 默认映射使用索引键范围Azure Blob 与有序的通用 S3 列表可定位下界并在上界停止S3 Express 与 Redis 在 Provider 定义遍历中过滤时间边界。恢复延迟包含所有选中 Catalog 的列表请求与存储客户端重试延迟因此应根据恢复延迟要求配置存储客户端的请求超时与重试上限。调优直觉README 明确给出更短的前瞻周期减少恢复分片的过早加载更短的检查间隔提高扫描频率减少新插入/新合格分片的等待时间。公开的JobShardManager.AssignJobShardsAsync方法会把相同的有序发现流收集进完整结果列表。优雅关闭生命周期Shutdown Lifecycle调度、取消请求与激活共享Orleans.Core中的Orleans.Internal.AdmissionGate工具做无锁准入控制admission调用方把返回的只读 token 保存在using局部变量中开始工作前检查Entereddispose 时释放准入每个准入 token 只有一个所有者、恰好 dispose 一次原子自增先占名额并同时观察关闭标志观察到关闭的尝试立即释放名额初步检查会让在自增前就看到关闭的调用者被拒绝因此后续到达不会改变 drain 计数调度与取消请求在整个完成过程中持有准入包括所有权查找与远端取消路由激活持有准入直到其执行任务被发布并入队。关闭时CloseAsync原子设置关闭标志、向进行中的请求发出取消信号并在快照运行中的分片并取消执行之前等待这些操作。回调失败会被记录而关闭继续 drain 请求、等待执行并释放分片。调度与取消的成功写入保留结果即使取消与完成竞态成功创建的分片仍保持已拥有。关闭随后等待当前扫描以及每个已准入分片的执行与清理用关闭 token 注销仍处于非激活状态的缓存分片并 dispose。Journaled Provider 会把有内容的分片释放给其他 Silo 认领并删除空分片包括请求取消后完成的创建。使用示例1. 实现IDurableJobHandler接口预约通知using Orleans; using Orleans.DurableJobs; public interface INotificationGrain : IGrainWithStringKey { Task ScheduleNotification(string message, DateTimeOffset sendTime); Task CancelScheduledNotification(CancellationToken requestCancellationToken); } public class NotificationGrain : Grain, INotificationGrain, IDurableJobHandler { private readonly ILocalDurableJobManager _jobManager; private readonly ILoggerNotificationGrain _logger; private DurableJob? _durableJob; public NotificationGrain( ILocalDurableJobManager jobManager, ILoggerNotificationGrain logger) { _jobManager jobManager; _logger logger; } public async Task ScheduleNotification(string message, DateTimeOffset sendTime) { var userId this.GetPrimaryKeyString(); var metadata new Dictionarystring, string { [Message] message }; _durableJob await _jobManager.ScheduleJobAsync( new ScheduleJobRequest { Target this.GetGrainId(), JobName SendNotification, DueTime sendTime, Metadata metadata }, CancellationToken.None); _logger.LogInformation( Scheduled notification for user {UserId} at {SendTime} (JobId: {JobId}), userId, sendTime, _durableJob.Id); } public async Task CancelScheduledNotification(CancellationToken requestCancellationToken) { if (_durableJob is null) { _logger.LogWarning(No scheduled notification to cancel); return; } var cancellationRequested await _jobManager.CancelAsync(_durableJob, requestCancellationToken); _logger.LogInformation( Notification {JobId} cancellation request recorded: {CancellationRequested}, _durableJob.Id, cancellationRequested); if (cancellationRequested) { // No future attempt will start. An already-running attempt may still complete. _durableJob null; } } // This method is called when the durable job executes public Task ExecuteJobAsync(IJobRunContext context, CancellationToken attemptCancellationToken) { var userId this.GetPrimaryKeyString(); var message context.Job.Metadata?[Message]; _logger.LogInformation( Sending notification to user {UserId}: {Message} (Job: {JobId}, Run: {RunId}, Attempt: {DequeueCount}), userId, message, context.Job.Id, context.RunId, context.DequeueCount); // Send the notification here // If this throws an exception, the job can be retried based on your retry policy _durableJob null; return Task.CompletedTask; } }2. 多任务工作流订单流程public interface IOrderGrain : IGrainWithGuidKey { Task PlaceOrder(OrderDetails details); Task CancelOrder(); } public class OrderGrain : Grain, IOrderGrain, IDurableJobHandler { private readonly ILocalDurableJobManager _jobManager; private readonly IOrderService _orderService; private readonly IGrainFactory _grainFactory; private readonly ILoggerOrderGrain _logger; public OrderGrain( ILocalDurableJobManager jobManager, IOrderService orderService, IGrainFactory grainFactory, ILoggerOrderGrain logger) { _jobManager jobManager; _orderService orderService; _grainFactory grainFactory; _logger logger; } public async Task PlaceOrder(OrderDetails details) { var orderId this.GetPrimaryKey(); // Create the order await _orderService.CreateOrderAsync(orderId, details); // Schedule delivery reminder for 24 hours before delivery var reminderTime details.DeliveryDate.AddHours(-24); await _jobManager.ScheduleJobAsync( new ScheduleJobRequest { Target this.GetGrainId(), JobName DeliveryReminder, DueTime reminderTime, Metadata new Dictionarystring, string { [Step] DeliveryReminder, [CustomerId] details.CustomerId, [OrderNumber] details.OrderNumber } }, CancellationToken.None); // Schedule order expiration if payment not received var expirationTime DateTimeOffset.UtcNow.AddHours(24); await _jobManager.ScheduleJobAsync( new ScheduleJobRequest { Target this.GetGrainId(), JobName OrderExpiration, DueTime expirationTime, Metadata new Dictionarystring, string { [Step] OrderExpiration } }, CancellationToken.None); } public async Task CancelOrder() { var orderId this.GetPrimaryKey(); await _orderService.CancelOrderAsync(orderId); } public async Task ExecuteJobAsync(IJobRunContext context, CancellationToken attemptCancellationToken) { var step context.Job.Metadata![Step]; var orderId this.GetPrimaryKey(); switch (step) { case DeliveryReminder: await HandleDeliveryReminder(context, attemptCancellationToken); break; case OrderExpiration: await HandleOrderExpiration(attemptCancellationToken); break; } } private async Task HandleDeliveryReminder(IJobRunContext context, CancellationToken attemptCancellationToken) { var customerId context.Job.Metadata![CustomerId]; var orderNumber context.Job.Metadata[OrderNumber]; var notificationGrain _grainFactory.GetGrainINotificationGrain(customerId); await notificationGrain.ScheduleNotification( $Your order #{orderNumber} will be delivered tomorrow!, DateTimeOffset.UtcNow); } private async Task HandleOrderExpiration(CancellationToken attemptCancellationToken) { var orderId this.GetPrimaryKey(); var order await _orderService.GetOrderAsync(orderId, attemptCancellationToken); if (order?.Status OrderStatus.Pending) { await _orderService.CancelOrderAsync(orderId, attemptCancellationToken); _logger.LogInformation(Order {OrderId} expired and canceled, orderId); } } }3. 带重试逻辑的任务public class PaymentProcessorGrain : Grain, IDurableJobHandler { private readonly IPaymentService _paymentService; private readonly ILoggerPaymentProcessorGrain _logger; public Task ExecuteJobAsync(IJobRunContext context, CancellationToken attemptCancellationToken) { var paymentId context.Job.Metadata?[PaymentId]; _logger.LogInformation( Processing payment {PaymentId} (Attempt {Attempt}), paymentId, context.DequeueCount); try { await _paymentService.ProcessPaymentAsync(paymentId, attemptCancellationToken); return Task.CompletedTask; } catch (TransientException ex) { _logger.LogWarning(ex, Payment processing failed with transient error, will retry); throw; // Let the retry policy handle it } catch (Exception ex) { _logger.LogError(ex, Payment processing failed with permanent error); throw; // This will not be retried if the retry policy returns null } } }4. 跟踪任务完成public class WorkflowGrain : Grain, IDurableJobHandler { private readonly Dictionarystring, TaskCompletionSource _pendingJobs new(); public async TaskDurableJob ScheduleWorkflowStep(string stepName, DateTimeOffset executeAt) { var job await _jobManager.ScheduleJobAsync( new ScheduleJobRequest { Target this.GetGrainId(), JobName stepName, DueTime executeAt, Metadata null }, CancellationToken.None); _pendingJobs[job.Id] new TaskCompletionSource(); return job; } public async Task WaitForJobCompletion(string jobId, TimeSpan timeout) { if (_pendingJobs.TryGetValue(jobId, out var tcs)) { using var cts new CancellationTokenSource(timeout); await tcs.Task.WaitAsync(cts.Token); } } public Task ExecuteJobAsync(IJobRunContext context, CancellationToken attemptCancellationToken) { // Execute the workflow step... // Mark as complete if (_pendingJobs.TryRemove(context.Job.Id, out var tcs)) { tcs.SetResult(); } return Task.CompletedTask; } }工作原理架构概览任务分片Job Sharding任务按时间划分到分片中分片默认时长由ShardDuration决定源码默认值为1 小时README 旧表述中的 1-minute 与当前源码不一致请以DurableJobsOptions.ShardDuration TimeSpan.FromHours(1)为准分片所有权Shard Ownership每个分片由单个 Silo 拥有并执行自动再平衡Silo 故障时其分片自动重新分配给健康 Silo有序执行分片内任务按到期时间顺序处理并发控制MaxConcurrentJobsPerSilo限制并发执行的任务数。任务生命周期┌─────────────┐ │ Scheduled │ ──▶ Job is created and added to appropriate shard └─────────────┘ │ ▼ ┌─────────────┐ │ Waiting │ ──▶ Job waits in queue until due time └─────────────┘ │ ▼ ┌─────────────┐ │ Executing │ ──▶ Job handler is invoked on target grain └─────────────┘ │ ├──▶ Success ──▶ Job is removed │ └──▶ Failure ──▶ Retry policy decides: • Retry: Job is re-queued with new due time • No Retry: Job is removed调度与执行的核心类型源码级ScheduleJobRequestScheduleJobRequest.cs必需字段Target接收任务的目标 GrainId、JobName非空任务名用于标识与 handler 路由、DueTime执行时间可选MetadataIReadOnlyDictionarystring, string以及 W3CTraceParent/TraceState——未显式指定时会自动捕获Activity.Current的 trace 上下文让任务执行时延续分布式追踪。DurableJobDurableJob.cs携带Id、Name、DueTime、TargetGrainId、ShardId、Metadata与 trace 字段是可序列化句柄。IJobRunContextIDurableJobHandler.cs提供Job、RunId每次出队生成一次重试/分片重分配后不保留与DequeueCount包括重试在内的出队次数。ExecuteJobAsync收到的attemptCancellationToken只取消当前这次尝试任务本身仍可被另一台主机再次尝试——这与CancelAsync的持久取消语义相互配合。运行结果DurableJobRunResult.cshandler 可通过返回Completed、InProgress(delay)任务仍在进行执行器按 delay 轮询期间继续占用并发槽位、Failed(exception)交给重试策略或RescheduleAt(dueTime)成功后持久重调度并重置失败尝试计数来精确控制后续行为。其中RescheduleRequested序列化值为 3滚动升级期间应在所有执行器都支持它之后再使用。配置参考DurableJobsOptions 完整字段下表基于当前源码 DurableJobsOptions.cs 整理与 README 的简易表格相比补齐了默认值与约束请以源码为准属性类型默认值说明ShardDurationTimeSpan1 小时每个任务分片的时长。越小延迟越低但开销越大。为与整点对齐建议选能整除 60 分钟的值1、2、3、4、5、6、10、12、15、20、30、60 分钟避免跨小时桶漂移。必须 0。ShardActivationBufferPeriodTimeSpan5 分钟分片在开始时间之前多久开始处理避免长时间持有空闲分片。ShardLoadLookaheadPeriodTimeSpan10 分钟提前多久加载合格分片。零表示选择开始时间不早于当前时间的分片发现视界上限为DateTimeOffset.MaxValue必须非负。ShardCheckIntervalTimeSpan5 分钟周期性发现与可写分片清理间隔成员变化也触发全新扫描。必须在 14294967294 毫秒之间。ShardStripeCountint1每个时间桶使用的可写分片数。增大可将同一到期桶的任务分散到多个分片 Journal必须在 132768 之间。分片选择是写侧的轮询扇出旋钮。JobStatusPollIntervalTimeSpan1 秒轮询异步任务 handler 前的等待时长轮询期间任务继续占用并发槽位必须 0。MaxConcurrentJobsPerSiloint10000 × CPU 核数单 Silo 可并发执行的最大任务数。OverloadBackoffDelayTimeSpan5 秒过载时任务批处理暂停的时长之后重新检查过载状态。ConcurrencySlowStartEnabledbooltrue是否启用并发慢启动从SlowStartInitialConcurrency开始每SlowStartInterval翻倍直到达到上限避免缓存、连接池、线程池预热前的饿死问题。SlowStartInitialConcurrencyintCPU 核数慢启动的初始并发必须 0。SlowStartIntervalTimeSpan10 秒慢启动并发翻倍间隔启用慢启动时必须 0。ShouldRetryFuncIJobRunContext, Exception, DateTimeOffset?最多 5 次指数退避2^n 秒失败任务是否重试及重试时间返回新到期时间或 null不重试不可为 null。ShardBatchLingerDelayTimeSpanZero分片操作处理器等待后续变更加入同一批 Journal 写入的时长。正数会以首个请求的有界延迟换取突发负载下的更大批量Zero 表示行为不变。必须非负。MaxAdoptedCountint3分片从死亡所有者处被收养的次数上限超过则标记为中毒poisoned不再分配给任何 Silo。仅从崩溃 Silo 收养时递增优雅关闭释放所有权不递增分片全部任务成功处理后被收养计数重置为 0。必须 ≥ 0。ShardClaimInitialBudgetint2启动后立即允许认领的孤儿分片数随后线性增长到ShardClaimMaxBudget。防止灾后恢复时新 Silo 一次性认领过多。必须 ≥ 0。ShardClaimMaxBudgetint20认领斜坡期结束时可累计认领的孤儿分片总数必须 ≥ShardClaimInitialBudget。ShardClaimRampUpDurationTimeSpan5 分钟认领斜坡期时长过后不再限制认领。设为Zero完全禁用斜坡。必须非负。ActiveProviderNamestringDefault用于创建新任务分片的活跃 Journal Provider 名。DrainingProviderNamesListstring空额外命名的 Journal Provider其已有分片被发现并 drain。已有任务保留原 Provider 处理重试、取消与清理。Provider 选择在 Silo 启动时固定。校验规则由 DurableJobsOptions.cs 中的DurableJobsOptionsValidator在 Silo 配置期强制执行如ShardDuration 0、ShardCheckInterval范围、ShardStripeCount范围、慢启动参数联动、认领预算单调性等违反会抛出OrleansConfigurationException同时DurableJobsJournalingConfigurationValidator会确认注册的是JournaledJobShardManager否则提示使用UseJournaledDurableJobs(...)配置存储。最佳实践设置合理的并发上限防止资源耗尽options.MaxConcurrentJobsPerSilo 100; // Adjust based on your workload实现幂等的任务 handler任务可能被重试确保 handler 幂等public async Task ExecuteJobAsync(IJobRunContext context, CancellationToken attemptCancellationToken) { var jobId context.Job.Id; // Check if already processed if (await _state.IsProcessed(jobId)) return; // Process job... await _state.MarkProcessed(jobId); }明智使用元数据保持轻量// Good: Store IDs var metadata new Dictionarystring, string { [OrderId] 12345 }; // Bad: Store large objects var metadata new Dictionarystring, string { [Order] JsonSerializer.Serialize(largeOrder) };处理取消尊重attemptCancellationTokenpublic async Task ExecuteJobAsync(IJobRunContext context, CancellationToken attemptCancellationToken) { await SomeLongRunningOperation(attemptCancellationToken); }关注持久化语义开发/测试用UseInMemoryDurableJobs进程内易失存储生产必须使用UseAzureBlobDurableJobs、UseAzureTableDurableJobs或命名 Redis/S3/Volatile Provider以支撑重启恢复。为迁移预留权限与观测A→B 迁移中保留 A 的列出/读取/认领/更新/重试/取消/压缩/删除权限并通过IDurableJobsStorageInspector持续盘点直到完全清零后再退役 A。延伸资料核心包源码Orleans.DurableJobs、DurableJobsExtensions.csAzure 存储后端Azure Storage 版 README迁移示例samples/DurableJobsMigration/README.md迁移部署与监控指南durable-jobs-migration.md可运行实验场Blob/Table 后端切换DurableJobsJournaling playground测试覆盖Durable Jobs 行为测试位于 test/Orleans.DurableJobs.Tests赞分享后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载相关推荐SLIM与Kubernetes Jobs一次性任务优化SLIM与Kubernetes Jobs一次性任务优化 引言Kubernetes Jobs的隐形痛点 你是否曾遇到Kubernetes Jobs执行缓慢、资云原生CLI应用安全终极指南使用Orleans Scheduled Jobs与Reminders轻松解决分布式定时任务难题终极指南使用Orleans Scheduled Jobs与Reminders轻松解决分布式定时任务难题 在现代分布式系统中定时任务和提醒功能是构建可靠服务的后端微服务Orleans Durable Jobs 命名 Journal Provider 与存储迁移指南单写 Provider Draining Providers 的滚动切换方案Orleans Durable Jobs 命名 Journal Provider 与存储迁移指南单写 Provider Draining Provider后端微服务上一篇GitHub中文化插件3分钟让你的GitHub界面全面变中文下一篇5分钟彻底告别GitHub英文界面中文翻译插件让你的开发效率飙升300%创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考