58
func NewArchivalQueueFactory(
59
params ArchivalQueueFactoryParams,
61
>
return &archivalQueueFactory{
62
>
ArchivalQueueFactoryParams: params,
63
>
QueueFactoryBase: newQueueFactoryBase(params),
64
>
}
65
>
}
66
67
// newHostScheduler creates a new task scheduler for tasks on the archival queue.
69
>
return queues.NewScheduler(
70
>
params.ClusterMetadata.GetCurrentClusterName(),
71
>
queues.SchedulerOptions{
72
>
WorkerCount: params.Config.ArchivalProcessorSchedulerWorkerCount,
73
>
ActiveNamespaceWeights: dynamicconfig.GetMapPropertyFnFilteredByNamespace(ArchivalTaskPriorities),
74
>
StandbyNamespaceWeights: dynamicconfig.GetMapPropertyFnFilteredByNamespace(ArchivalTaskPriorities),
75
>
InactiveNamespaceDeletionDelay: params.Config.TaskSchedulerInactiveChannelDeletionDelay,
76
>
ExecutionAwareSchedulerOptions: ctasks.ExecutionAwareSchedulerOptions{
77
>
Enabled: params.Config.TaskSchedulerEnableExecutionQueueScheduler,
78
>
MaxQueues: params.Config.TaskSchedulerExecutionQueueSchedulerMaxQueues,
79
>
QueueTTL: params.Config.TaskSchedulerExecutionQueueSchedulerQueueTTL,
80
>
QueueConcurrency: params.Config.TaskSchedulerExecutionQueueSchedulerQueueConcurrency,
81
>
},
82
>
},
83
>
params.NamespaceRegistry,
84
>
params.Logger,
85
>
params.MetricsHandler,
86
>
params.TimeSource,
87
>
)
88
>
}
89
90
// newQueueFactoryBase creates a new QueueFactoryBase for the archival queue, which contains common configurations
91
// like the task scheduler, task priority assigner, and rate limiters.
93
>
return QueueFactoryBase{
94
>
HostScheduler: newHostScheduler(params),
95
>
HostPriorityAssigner: queues.NewPriorityAssigner(
96
>
params.NamespaceRegistry,
97
>
params.ClusterMetadata.GetCurrentClusterName(),
98
>
),
99
>
HostReaderRateLimiter: queues.NewReaderPriorityRateLimiter(
100
>
NewHostRateLimiterRateFn(
101
>
params.Config.ArchivalProcessorMaxPollHostRPS,
102
>
params.Config.PersistenceMaxQPS,
103
>
archivalQueuePersistenceMaxRPSRatio,
104
>
),
105
>
int64(params.Config.ArchivalQueueMaxReaderCount()),
106
>
),
107
>
Tracer: params.TracerProvider.Tracer(telemetry.ComponentQueueArchival),
108
>
}
109
>
}
110
111
// CreateQueue creates a new archival queue for the given shard.