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

Method removePartitions

offset_manager.go:501–532  ·  view source on GitHub ↗

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)

Source from the content-addressed store, hash-verified

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)
501func (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)
535func (om *offsetManager) releasePOMs(force bool) (remaining int) {

Callers 1

Calls 3

releaseSelectedPOMsMethod · 0.95
flushToBrokerForMethod · 0.95
AsyncCloseMethod · 0.65

Tested by 1