mirror of
https://github.com/minio/minio.git
synced 2025-01-24 05:03:16 -05:00
replication: queue existing objects to same workers as incoming (#18020)
Previously existing objects were queued to single worker and MRF re-queues are also handled by same worker - this does not fully use the available bandwidth in case there is no incoming workload.
This commit is contained in:
parent
c8a57a8fa2
commit
96fbf18201
@ -1661,9 +1661,8 @@ type ReplicationPool struct {
|
||||
resyncer *replicationResyncer
|
||||
|
||||
// workers:
|
||||
workers []chan ReplicationWorkerOperation
|
||||
lrgworkers []chan ReplicationWorkerOperation
|
||||
existingWorkers chan ReplicationWorkerOperation
|
||||
workers []chan ReplicationWorkerOperation
|
||||
lrgworkers []chan ReplicationWorkerOperation
|
||||
|
||||
// mrf:
|
||||
mrfWorkerKillCh chan struct{}
|
||||
@ -1723,8 +1722,6 @@ func NewReplicationPool(ctx context.Context, o ObjectLayer, opts replicationPool
|
||||
pool := &ReplicationPool{
|
||||
workers: make([]chan ReplicationWorkerOperation, 0, workers),
|
||||
lrgworkers: make([]chan ReplicationWorkerOperation, 0, LargeWorkerCount),
|
||||
existingWorkers: make(chan ReplicationWorkerOperation, 100000),
|
||||
|
||||
mrfReplicaCh: make(chan ReplicationWorkerOperation, 100000),
|
||||
mrfWorkerKillCh: make(chan struct{}, failedWorkers),
|
||||
resyncer: newresyncer(),
|
||||
@ -1738,7 +1735,6 @@ func NewReplicationPool(ctx context.Context, o ObjectLayer, opts replicationPool
|
||||
pool.AddLargeWorkers()
|
||||
pool.ResizeWorkers(workers, 0)
|
||||
pool.ResizeFailedWorkers(failedWorkers)
|
||||
go pool.AddWorker(pool.existingWorkers, nil)
|
||||
go pool.resyncer.PersistToDisk(ctx, o)
|
||||
go pool.processMRF()
|
||||
go pool.persistMRF()
|
||||
@ -1965,9 +1961,7 @@ func (p *ReplicationPool) queueReplicaTask(ri ReplicateObjectInfo) {
|
||||
|
||||
var ch, healCh chan<- ReplicationWorkerOperation
|
||||
switch ri.OpType {
|
||||
case replication.ExistingObjectReplicationType:
|
||||
ch = p.existingWorkers
|
||||
case replication.HealReplicationType:
|
||||
case replication.HealReplicationType, replication.ExistingObjectReplicationType:
|
||||
ch = p.mrfReplicaCh
|
||||
healCh = p.getWorkerCh(ri.Name, ri.Bucket, ri.Size)
|
||||
default:
|
||||
@ -2026,9 +2020,7 @@ func (p *ReplicationPool) queueReplicaDeleteTask(doi DeletedObjectReplicationInf
|
||||
}
|
||||
var ch chan<- ReplicationWorkerOperation
|
||||
switch doi.OpType {
|
||||
case replication.ExistingObjectReplicationType:
|
||||
ch = p.existingWorkers
|
||||
case replication.HealReplicationType:
|
||||
case replication.HealReplicationType, replication.ExistingObjectReplicationType:
|
||||
fallthrough
|
||||
default:
|
||||
ch = p.getWorkerCh(doi.Bucket, doi.ObjectName, 0)
|
||||
|
Loading…
x
Reference in New Issue
Block a user