125
Logger: params.Logger,
126
KeyFn: func(e queues.Executable) tasks.TaskGroupNamespaceIDAndDestination {
128
>
},
129
>
RunnableFactory: func(e queues.Executable) ctasks.Runnable {
130
>
key := grouper.KeyTyped(e.GetTask())
131
>
nsName := getNamespaceNameOrDefault(
132
>
params.NamespaceRegistry,
133
>
key.NamespaceID,
134
>
key.NamespaceID,
135
>
metricsHandler,
136
>
)
137
>
taggedMetricsHandler := metricsHandler.WithTags(
138
>
metrics.NamespaceTag(nsName),
139
>
metrics.DestinationTag(key.Destination),
140
>
)
141
>
return ctasks.NewRateLimitedTaskRunnableFromTask(
142
>
ctasks.RunnableTask{
143
>
Task: queues.NewCircuitBreakerExecutable(
144
>
e,
145
>
params.CircuitBreakerPool.Get(key),
146
>
taggedMetricsHandler,
147
>
),
148
>
},
149
>
rateLimiterPool.Get(key),
150
>
taggedMetricsHandler,
151
>
)
152
>
},
153
SchedulerFactory: func(
154
key tasks.TaskGroupNamespaceIDAndDestination,
156
>
nsName := getNamespaceNameOrDefault(
157
>
params.NamespaceRegistry,
158
>
key.NamespaceID,
159
>
key.NamespaceID,
160
>
metricsHandler,
161
>
)
162
>
return ctasks.NewDynamicWorkerPoolScheduler(
163
>
groupLimiter{
164
>
key: key,
165
>
namespaceRegistry: params.NamespaceRegistry,
166
>
metricsHandler: metricsHandler,
167
>
bufferSize: params.Config.OutboundQueueGroupLimiterBufferSize,
168
>
concurrency: params.Config.OutboundQueueGroupLimiterConcurrency,
169
>
},
170
>
metricsHandler.WithTags(
171
>
metrics.NamespaceTag(nsName),
172
>
metrics.DestinationTag(key.Destination),
173
>
),
174
>
)
175
>
},
176
},
177
),