|
@@ -462,11 +462,19 @@ func (bom *brokerOffsetManager) flushToBroker() {
|
|
|
case ErrNoError:
|
|
|
block := request.blocks[s.topic][s.partition]
|
|
|
s.updateCommitted(block.offset, block.metadata)
|
|
|
- break
|
|
|
- case ErrUnknownTopicOrPartition, ErrNotLeaderForPartition, ErrLeaderNotAvailable:
|
|
|
+ case ErrUnknownTopicOrPartition, ErrNotLeaderForPartition, ErrLeaderNotAvailable,
|
|
|
+ ErrConsumerCoordinatorNotAvailable, ErrNotCoordinatorForConsumer:
|
|
|
+
|
|
|
delete(bom.subscriptions, s)
|
|
|
s.rebalance <- none{}
|
|
|
+ case ErrOffsetMetadataTooLarge, ErrInvalidCommitOffsetSize:
|
|
|
+
|
|
|
+ s.handleError(err)
|
|
|
+ case ErrOffsetsLoadInProgress:
|
|
|
+
|
|
|
+ break
|
|
|
default:
|
|
|
+
|
|
|
s.handleError(err)
|
|
|
delete(bom.subscriptions, s)
|
|
|
s.rebalance <- none{}
|