112
func (f *archivalQueueFactory) CreateQueue(
113
shard historyi.ShardContext,
115
>
executor := f.newArchivalTaskExecutor(shard, f.WorkflowCache)
116
>
if f.ExecutorWrapper != nil {
117
executor = f.ExecutorWrapper.Wrap(executor)
118
}
120
}
121
122
// newArchivalTaskExecutor creates a new archival task executor for the given shard.
123
>
func (f *archivalQueueFactory) newArchivalTaskExecutor(shard historyi.ShardContext, workflowCache wcache.Cache) queues.Executor {
archival_queue_factory.go
124
>
return NewArchivalQueueTaskExecutor(
125
>
f.Archiver,
126
>
shard,
127
>
workflowCache,
128
>
f.RelocatableAttributesFetcher,
129
>
f.MetricsHandler,
130
>
log.With(shard.GetLogger(), tag.ComponentArchivalQueue),
131
>
)
132
>
}
133
134
// newScheduledQueue creates a new scheduled queue for the given shard with archival-specific configurations.
135
>
func (f *archivalQueueFactory) newScheduledQueue(shard historyi.ShardContext, executor queues.Executor) queues.Queue {
archival_queue_factory.go
136
>
logger := log.With(shard.GetLogger(), tag.ComponentArchivalQueue)
137
>
metricsHandler := f.MetricsHandler.WithTags(metrics.OperationTag(metrics.OperationArchivalQueueProcessorScope))
138
>
139
>
shardScheduler := queues.NewRateLimitedScheduler(
140
>
f.HostScheduler,
141
>
queues.RateLimitedSchedulerOptions{
142
>
Enabled: f.Config.TaskSchedulerEnableRateLimiter,
143
>
EnableShadowMode: f.Config.TaskSchedulerEnableRateLimiterShadowMode,
144
>
StartupDelay: f.Config.TaskSchedulerRateLimiterStartupDelay,
145
>
},
146
>
f.ClusterMetadata.GetCurrentClusterName(),
147
>
f.NamespaceRegistry,
148
>
f.SchedulerRateLimiter,
149
>
f.TimeSource,
150
>
f.ChasmRegistry,
151
>
logger,
152
>
metricsHandler,
153
>
)
154
>
155
>
rescheduler := queues.NewRescheduler(
156
>
shardScheduler,
157
>
shard.GetTimeSource(),
158
>
logger,
159
>
metricsHandler,
160
>
)
161
>
162
>
factory := queues.NewExecutableFactory(
163
>
executor,
164
>
shardScheduler,
165
>
rescheduler,
166
>
f.HostPriorityAssigner,
167
>
shard.GetTimeSource(),
168
>
shard.GetNamespaceRegistry(),
169
>
shard.GetClusterMetadata(),
170
>
f.ChasmRegistry,
171
>
queues.GetTaskTypeTagValue,
172
>
logger,
173
>
metricsHandler,
174
>
f.Tracer,
175
>
f.DLQWriter,
176
>
f.Config.TaskDLQEnabled,
177
>
f.Config.TaskDLQUnexpectedErrorAttempts,
178
>
f.Config.TaskDLQInternalErrors,
179
>
f.Config.TaskDLQErrorPattern,
180
>
)
181
>
return queues.NewScheduledQueue(
182
>
shard,
183
>
tasks.CategoryArchival,
184
>
shardScheduler,
185
>
rescheduler,
186
>
factory,
187
>
&queues.Options{
188
>
ReaderOptions: queues.ReaderOptions{
189
>
BatchSize: f.Config.ArchivalTaskBatchSize,
190
>
MaxPendingTasksCount: f.Config.QueuePendingTaskMaxCount,
191
>
PollBackoffInterval: f.Config.ArchivalProcessorPollBackoffInterval,
192
>
MaxPredicateSize: f.Config.QueueMaxPredicateSize,
193
>
},
194
>
MonitorOptions: queues.MonitorOptions{
195
>
PendingTasksCriticalCount: f.Config.QueuePendingTaskCriticalCount,
196
>
ReaderStuckCriticalAttempts: f.Config.QueueReaderStuckCriticalAttempts,
197
>
SliceCountCriticalThreshold: f.Config.QueueCriticalSlicesCount,
198
>
},
199
>
MaxPollRPS: f.Config.ArchivalProcessorMaxPollRPS,
200
>
MaxPollInterval: f.Config.ArchivalProcessorMaxPollInterval,
201
>
MaxPollIntervalJitterCoefficient: f.Config.ArchivalProcessorMaxPollIntervalJitterCoefficient,
202
>
CheckpointInterval: f.Config.ArchivalProcessorUpdateAckInterval,
203
>
CheckpointIntervalJitterCoefficient: f.Config.ArchivalProcessorUpdateAckIntervalJitterCoefficient,
204
>
MaxReaderCount: f.Config.ArchivalQueueMaxReaderCount,
205
>
MoveGroupTaskCountBase: f.Config.QueueMoveGroupTaskCountBase,
206
>
MoveGroupTaskCountMultiplier: f.Config.QueueMoveGroupTaskCountMultiplier,
207
>
ShrinkPredicateMaxPendingKeys: f.Config.QueueShrinkPredicateMaxPendingKeys,
208
>
},
209
>
f.HostReaderRateLimiter,
210
>
logger,
211
>
metricsHandler,
212
>
)
213
>
}