Merge pull request #17 from vmware/spentakota/sendLeaseRenewedMetric
feat: Sending renewed lease metric
This commit is contained in:
commit
c5bc6c4ded
2 changed files with 4 additions and 0 deletions
|
|
@ -103,6 +103,8 @@ func (sc *FanOutShardConsumer) getRecords() error {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
refreshLeaseTimer = time.After(time.Until(sc.shard.LeaseTimeout.Add(-time.Duration(sc.kclConfig.LeaseRefreshPeriodMillis) * time.Millisecond)))
|
refreshLeaseTimer = time.After(time.Until(sc.shard.LeaseTimeout.Add(-time.Duration(sc.kclConfig.LeaseRefreshPeriodMillis) * time.Millisecond)))
|
||||||
|
// log metric for renewed lease for worker
|
||||||
|
sc.mService.LeaseRenewed(sc.shard.ID)
|
||||||
case event, ok := <-shardSub.GetStream().Events():
|
case event, ok := <-shardSub.GetStream().Events():
|
||||||
if !ok {
|
if !ok {
|
||||||
// need to resubscribe to shard
|
// need to resubscribe to shard
|
||||||
|
|
|
||||||
|
|
@ -122,6 +122,8 @@ func (sc *PollingShardConsumer) getRecords() error {
|
||||||
sc.shard.ID, sc.consumerID, err)
|
sc.shard.ID, sc.consumerID, err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
// log metric for renewed lease for worker
|
||||||
|
sc.mService.LeaseRenewed(sc.shard.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
getRecordsStartTime := time.Now()
|
getRecordsStartTime := time.Now()
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue