removePartitions closes partitions revoked during a rebalance, committing first when configured lint:ignore U1000 // consumed by the cooperative rebalancing path added in a following PR; the only in-build caller here is the unit test (excluded by the integration build tag)
(topicPartitions map[string][]int32)
| 499 | // |
| 500 | //lint:ignore U1000 // consumed by the cooperative rebalancing path added in a following PR; the only in-build caller here is the unit test (excluded by the integration build tag) |
| 501 | func (om *offsetManager) removePartitions(topicPartitions map[string][]int32) { |
| 502 | targets := make(partitionTargets, len(topicPartitions)) |
| 503 | for topic, partitions := range topicPartitions { |
| 504 | targets[topic] = make(map[int32]none, len(partitions)) |
| 505 | for _, partition := range partitions { |
| 506 | targets[topic][partition] = none{} |
| 507 | } |
| 508 | } |
| 509 | |
| 510 | om.pomsLock.RLock() |
| 511 | for topic, partitions := range targets { |
| 512 | for partition := range partitions { |
| 513 | if pom := om.poms[topic][partition]; pom != nil { |
| 514 | pom.AsyncClose() |
| 515 | } |
| 516 | } |
| 517 | } |
| 518 | om.pomsLock.RUnlock() |
| 519 | |
| 520 | if !om.conf.Consumer.Offsets.AutoCommit.Enable { |
| 521 | om.releaseSelectedPOMs(true, targets) |
| 522 | return |
| 523 | } |
| 524 | |
| 525 | for attempt := 0; attempt <= om.conf.Consumer.Offsets.Retry.Max; attempt++ { |
| 526 | om.flushToBrokerFor(targets) |
| 527 | if om.releaseSelectedPOMs(false, targets) == 0 { |
| 528 | return |
| 529 | } |
| 530 | } |
| 531 | om.releaseSelectedPOMs(true, targets) |
| 532 | } |
| 533 | |
| 534 | // Releases/removes closed POMs once they are clean (or when forced) |
| 535 | func (om *offsetManager) releasePOMs(force bool) (remaining int) { |