193
config *configs.Config,
194
clientBean client.Bean,
196
>
numWorker := config.ReplicationTaskFetcherParallelism()
197
>
requestChan := make(chan *replicationTaskRequest, requestChanBufferSize)
198
>
shutdownChan := make(chan struct{})
199
>
rateLimiter := quotas.NewDefaultOutgoingRateLimiter(
200
>
func() float64 { return config.ReplicationTaskProcessorHostQPS() },
201
)
202
203
>
workers := make(map[int]*replicationTaskFetcherWorker)
task_fetcher.go
204
>
for i := range numWorker {
205
>
workers[i] = newReplicationTaskFetcherWorker(
206
>
logger,
207
>
sourceCluster,
208
>
currentCluster,
209
>
config,
210
>
clientBean,
211
>
rateLimiter,
212
>
requestChan,
213
>
shutdownChan,
214
>
)
215
>
}
216
218
>
status: common.DaemonStatusInitialized,
219
>
config: config,
220
>
numWorker: numWorker,
221
>
logger: log.With(logger, tag.ClusterName(sourceCluster)),
222
>
currentCluster: currentCluster,
223
>
sourceCluster: sourceCluster,
224
>
rateLimiter: rateLimiter,
225
>
requestChan: requestChan,
226
>
shutdownChan: shutdownChan,
227
>
workers: workers,
228
>
}
229
}
230