MCPcopy Create free account
hub / github.com/IBM/sarama / fetchInitialOffset

Method fetchInitialOffset

offset_manager.go:151–203  ·  view source on GitHub ↗
(topic string, partition int32, retries int)

Source from the content-addressed store, hash-verified

149}
150
151func (om *offsetManager) fetchInitialOffset(topic string, partition int32, retries int) (int64, int32, string, error) {
152 broker, err := om.coordinator()
153 if err != nil {
154 if retries <= 0 {
155 return 0, 0, "", err
156 }
157 return om.fetchInitialOffset(topic, partition, retries-1)
158 }
159
160 partitions := map[string][]int32{topic: {partition}}
161 req := NewOffsetFetchRequest(om.conf.Version, om.group, partitions)
162 resp, err := broker.FetchOffset(req)
163 if err != nil {
164 if retries <= 0 {
165 return 0, 0, "", err
166 }
167 om.releaseCoordinator(broker)
168 return om.fetchInitialOffset(topic, partition, retries-1)
169 }
170
171 block := resp.GetBlock(topic, partition)
172 if block == nil {
173 // v2+ surfaces some coordinator errors at the top level with no per-partition blocks
174 if resp.Err == ErrNoError {
175 return 0, 0, "", ErrIncompleteResponse
176 }
177 block = &OffsetFetchResponseBlock{Err: resp.Err}
178 }
179
180 switch block.Err {
181 case ErrNoError:
182 return block.Offset, block.LeaderEpoch, block.Metadata, nil
183 case ErrNotCoordinatorForConsumer, ErrConsumerCoordinatorNotAvailable:
184 if retries <= 0 {
185 return 0, 0, "", block.Err
186 }
187 om.releaseCoordinator(broker)
188 return om.fetchInitialOffset(topic, partition, retries-1)
189 case ErrOffsetsLoadInProgress:
190 if retries <= 0 {
191 return 0, 0, "", block.Err
192 }
193 backoff := om.computeBackoff(retries)
194 select {
195 case <-om.closing:
196 return 0, 0, "", block.Err
197 case <-time.After(backoff):
198 }
199 return om.fetchInitialOffset(topic, partition, retries-1)
200 default:
201 return 0, 0, "", block.Err
202 }
203}
204
205func (om *offsetManager) coordinator() (*Broker, error) {
206 om.brokerLock.RLock()

Callers 1

Calls 6

coordinatorMethod · 0.95
releaseCoordinatorMethod · 0.95
computeBackoffMethod · 0.95
NewOffsetFetchRequestFunction · 0.85
FetchOffsetMethod · 0.80
GetBlockMethod · 0.45

Tested by

no test coverage detected