(topic string, partition int32, retries int)
| 149 | } |
| 150 | |
| 151 | func (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 | |
| 205 | func (om *offsetManager) coordinator() (*Broker, error) { |
| 206 | om.brokerLock.RLock() |
no test coverage detected