103
logger log.Logger,
104
metricsHandler metrics.Handler,
106
>
107
>
sliceList := list.New()
108
>
for _, slice := range slices {
109
sliceList.PushBack(slice)
110
}
111
>
monitor.SetSliceCount(readerID, len(slices))
reader.go
112
>
113
>
rateLimitContext, rateLimitContextCancel := context.WithCancel(context.Background())
114
>
return &ReaderImpl{
115
>
readerID: readerID,
116
>
options: options,
117
>
scheduler: scheduler,
118
>
rescheduler: rescheduler,
119
>
timeSource: timeSource,
120
>
ratelimiter: ratelimiter,
121
>
monitor: monitor,
122
>
completionFn: completionFn,
123
>
logger: log.With(logger, tag.QueueReaderID(readerID)),
124
>
metricsHandler: metricsHandler,
125
>
126
>
status: common.DaemonStatusInitialized,
127
>
shutdownCh: make(chan struct{}),
128
>
129
>
slices: sliceList,
130
>
nextReadSlice: sliceList.Front(),
131
>
notifyCh: make(chan struct{}, 1),
132
>
133
>
retrier: backoff.NewRetrier(
134
>
common.CreateReadTaskRetryPolicy(),
135
>
clock.NewRealTimeSource(),
136
>
),
137
>
138
>
rateLimitContext: rateLimitContext,
139
>
rateLimitContextCancel: rateLimitContextCancel,
140
>
rateLimiterRequest: newReaderRequest(readerID),
141
>
}
142
}
143