85
subqueue subqueueIndex,
86
initialAckLevel fairLevel,
88
>
return &fairTaskReader{
89
>
backlogMgr: backlogMgr,
90
>
subqueue: subqueue,
91
>
logger: backlogMgr.logger,
92
>
retrier: backoff.NewRetrier(
93
>
backoff.NewExponentialRetryPolicy(50*time.Millisecond).
94
>
WithMaximumInterval(10*time.Second).
95
>
WithExpirationInterval(backoff.NoInterval),
96
>
clock.NewRealTimeSource(),
97
>
),
98
>
throttleRetrier: backoff.NewRetrier(
99
>
backoff.NewExponentialRetryPolicy(2*time.Second).
100
>
WithMaximumInterval(30*time.Second).
101
>
WithExpirationInterval(backoff.NoInterval),
102
>
clock.NewRealTimeSource(),
103
>
),
104
>
backlogAge: newBacklogAgeTracker(),
105
>
addRetries: semaphore.NewWeighted(concurrentAddRetries),
106
>
107
>
// ack manager
108
>
outstandingTasks: *newFairLevelTreeMap(),
109
>
readLevel: initialAckLevel,
110
>
ackLevel: initialAckLevel,
111
>
evictedAcks: *btree.NewBTreeGOptions(fairLevel.less, btree.Options{NoLocks: true}),
112
>
113
>
// gc state
114
>
lastGCTime: time.Now(),
115
>
}
116
>
}
117
119
>
tr.lock.Lock()
120
>
defer tr.lock.Unlock()
121
>
tr.maybeReadTasksLocked()
122
>
}
123
124
func (tr *fairTaskReader) getOldestBacklogTime() time.Time {