diff --git a/CHANGELOG.md b/CHANGELOG.md index bb87671c14f..9838bcaffe0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -91,6 +91,7 @@ To learn more about active deprecations, we recommend checking [GitHub Discussio - **General**: Restore gRPC reconnect backoff in the metrics service client; an unset `Backoff` in `WithConnectParams` disabled backoff and caused a zero-delay reconnect loop that flooded logs when keda-operator was unreachable ([#7856](https://github.com/kedacore/keda/issues/7856)) - **General**: Treat negative external metric values as zero to prevent incorrect HPA scaling ([#7880](https://github.com/kedacore/keda/issues/7880)) - **Azure Blob Storage Scaler**: Fix `globPattern` never matching when written in path-style with a leading `/`, since blob names never have one ([#6492](https://github.com/kedacore/keda/issues/6492)) +- **Kafka Scaler**: Fix overestimated lag for consumer groups without a committed offset when `offsetResetPolicy` is `earliest`; messages already deleted by retention are no longer counted ([#7936](https://github.com/kedacore/keda/issues/7936)) - **MongoDB Scaler**: Disconnect the client when the initial `Ping` fails so the background topology-monitoring connections opened by `mongo.Connect` are not leaked ([#5612](https://github.com/kedacore/keda/issues/5612)) ### Deprecations diff --git a/pkg/scalers/kafka_scaler.go b/pkg/scalers/kafka_scaler.go index f3aa42e73a1..d5195365c98 100644 --- a/pkg/scalers/kafka_scaler.go +++ b/pkg/scalers/kafka_scaler.go @@ -633,7 +633,7 @@ func (s *kafkaScaler) getConsumerOffsets(topicPartitions map[string][]int32) (*s // When excludePersistentLag is set to `false` (default), lag will always be equal to lagWithPersistent // When excludePersistentLag is set to `true`, if partition is deemed to have persistent lag, lag will be set to 0 and lagWithPersistent will be latestOffset - consumerOffset // These return values will allow proper scaling from 0 -> 1 replicas by the IsActive func. -func (s *kafkaScaler) getLagForPartition(topic string, partitionID int32, offsets *sarama.OffsetFetchResponse, topicPartitionOffsets map[string]map[int32]int64) (int64, int64, error) { +func (s *kafkaScaler) getLagForPartition(topic string, partitionID int32, offsets *sarama.OffsetFetchResponse, topicPartitionOffsets map[string]map[int32]partitionOffsets) (int64, int64, error) { block := offsets.GetBlock(topic, partitionID) if block == nil { errMsg := fmt.Errorf("error finding offset block for topic %s and partition %d from offset block: %v", topic, partitionID, offsets.Blocks) @@ -664,8 +664,8 @@ func (s *kafkaScaler) getLagForPartition(topic string, partitionID int32, offset return retVal, retVal, nil } // offsetResetPolicy == earliest - // For earliest policy, we need latestOffset to return the full lag when scaleToZeroOnInvalidOffset is false - // But if we can't get latestOffset, we should still respect scaleToZeroOnInvalidOffset + // For earliest policy, we need partition offsets to return the retained lag when scaleToZeroOnInvalidOffset is false + // But if we can't get partition offsets, we should still respect scaleToZeroOnInvalidOffset if s.metadata.ScaleToZeroOnInvalidOffset { return 0, 0, nil } @@ -676,18 +676,24 @@ func (s *kafkaScaler) getLagForPartition(topic string, partitionID int32, offset s.logger.V(1).Info(fmt.Sprintf("Topic %s not found in latest offset response, treating partition %d as 0 lag", topic, partitionID)) return 0, 0, nil } - latestOffset, partitionFound := topicOffsets[partitionID] + offsetsForPartition, partitionFound := topicOffsets[partitionID] if !partitionFound { // Partition missing from latest offset response - treat as 0 lag and continue // This can happen with Azure Event Hub when partitions are intermittently not returned s.logger.V(1).Info(fmt.Sprintf("Partition %d in topic %s not found in latest offset response, treating as 0 lag", partitionID, topic)) return 0, 0, nil } + latestOffset := offsetsForPartition.latestOffset // If we got here with invalidOffset and earliest policy, scaleToZeroOnInvalidOffset must be false - // Return the full lag (latestOffset) as per earliest policy behavior + // Return the retained lag from log start to log end if consumerOffset == invalidOffset && s.metadata.OffsetResetPolicy == earliest { - return latestOffset, latestOffset, nil + if !offsetsForPartition.earliestOffsetFound { + s.logger.V(1).Info(fmt.Sprintf("Partition %d in topic %s has no earliest offset, falling back to latest offset as lag", partitionID, topic)) + return latestOffset, latestOffset, nil + } + lag := latestOffset - offsetsForPartition.earliestOffset + return lag, lag, nil } // This code block tries to prevent KEDA Kafka trigger from scaling the scale target based on erroneous events @@ -758,12 +764,18 @@ type consumerOffsetResult struct { err error } +type partitionOffsets struct { + earliestOffset int64 + earliestOffsetFound bool + latestOffset int64 +} + type producerOffsetResult struct { producerOffsets map[string]map[int32]int64 err error } -func (s *kafkaScaler) getConsumerAndProducerOffsets(topicPartitions map[string][]int32) (*sarama.OffsetFetchResponse, map[string]map[int32]int64, error) { +func (s *kafkaScaler) getConsumerAndProducerOffsets(topicPartitions map[string][]int32) (*sarama.OffsetFetchResponse, map[string]map[int32]partitionOffsets, error) { consumerChan := make(chan consumerOffsetResult, 1) go func() { consumerOffsets, err := s.getConsumerOffsets(topicPartitions) @@ -772,7 +784,7 @@ func (s *kafkaScaler) getConsumerAndProducerOffsets(topicPartitions map[string][ producerChan := make(chan producerOffsetResult, 1) go func() { - producerOffsets, err := s.getProducerOffsets(topicPartitions) + producerOffsets, err := s.getProducerOffsets(topicPartitions, sarama.OffsetNewest) producerChan <- producerOffsetResult{producerOffsets, err} }() @@ -786,7 +798,44 @@ func (s *kafkaScaler) getConsumerAndProducerOffsets(topicPartitions map[string][ return nil, nil, producerRes.err } - return consumerRes.consumerOffsets, producerRes.producerOffsets, nil + topicPartitionOffsets := make(map[string]map[int32]partitionOffsets, len(producerRes.producerOffsets)) + for topic, partitionLatestOffsets := range producerRes.producerOffsets { + topicPartitionOffsets[topic] = make(map[int32]partitionOffsets, len(partitionLatestOffsets)) + for partition, latestOffset := range partitionLatestOffsets { + topicPartitionOffsets[topic][partition] = partitionOffsets{latestOffset: latestOffset} + } + } + + if s.metadata.OffsetResetPolicy == earliest && !s.metadata.ScaleToZeroOnInvalidOffset { + partitionsWithInvalidOffsets := make(map[string][]int32) + for topic, partitionOffsetsForTopic := range topicPartitionOffsets { + for partition := range partitionOffsetsForTopic { + consumerOffsetBlock := consumerRes.consumerOffsets.GetBlock(topic, partition) + if consumerOffsetBlock != nil && consumerOffsetBlock.Err == sarama.ErrNoError && consumerOffsetBlock.Offset == invalidOffset { + partitionsWithInvalidOffsets[topic] = append(partitionsWithInvalidOffsets[topic], partition) + } + } + } + + if len(partitionsWithInvalidOffsets) > 0 { + // Fetch earliest offsets only for partitions without a committed consumer offset. + // This avoids unnecessary broker requests when latest offsets are sufficient. + earliestOffsets, err := s.getProducerOffsets(partitionsWithInvalidOffsets, sarama.OffsetOldest) + if err != nil { + s.logger.Error(err, "error getting earliest offsets, falling back to latest offsets for lag calculation") + } + for topic, partitionEarliestOffsets := range earliestOffsets { + for partition, earliestOffset := range partitionEarliestOffsets { + offsetsForPartition := topicPartitionOffsets[topic][partition] + offsetsForPartition.earliestOffset = earliestOffset + offsetsForPartition.earliestOffsetFound = true + topicPartitionOffsets[topic][partition] = offsetsForPartition + } + } + } + } + + return consumerRes.consumerOffsets, topicPartitionOffsets, nil } // GetMetricsAndActivity returns value for a supported metric and an error if there is a problem getting the metric @@ -871,7 +920,7 @@ type brokerOffsetResult struct { err error } -func (s *kafkaScaler) getProducerOffsets(topicPartitions map[string][]int32) (map[string]map[int32]int64, error) { +func (s *kafkaScaler) getProducerOffsets(topicPartitions map[string][]int32, offsetPosition int64) (map[string]map[int32]int64, error) { version := int16(0) if s.client.Config().Version.IsAtLeast(sarama.V0_10_1_0) { version = 1 @@ -891,7 +940,7 @@ func (s *kafkaScaler) getProducerOffsets(topicPartitions map[string][]int32) (ma request = &sarama.OffsetRequest{Version: version} requests[broker] = request } - request.AddBlock(topic, partitionID, sarama.OffsetNewest, 1) + request.AddBlock(topic, partitionID, offsetPosition, 1) } } diff --git a/pkg/scalers/kafka_scaler_test.go b/pkg/scalers/kafka_scaler_test.go index 3ff057407f6..1f6696e5c4a 100644 --- a/pkg/scalers/kafka_scaler_test.go +++ b/pkg/scalers/kafka_scaler_test.go @@ -1018,7 +1018,7 @@ func TestGetLagForPartition_MissingPartition(t *testing.T) { tests := []struct { name string consumerOffset int64 - topicPartitionOffsets map[string]map[int32]int64 + topicPartitionOffsets map[string]map[int32]partitionOffsets offsetResetPolicy offsetResetPolicy scaleToZeroOnInvalidOffset bool expectedLag int64 @@ -1029,7 +1029,7 @@ func TestGetLagForPartition_MissingPartition(t *testing.T) { { name: "Scenario 1: scaleToZeroOnInvalidOffset true, invalid consumer offset, missing partition", consumerOffset: invalidOffset, - topicPartitionOffsets: map[string]map[int32]int64{ + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ "test-topic": { // Partition 0 is missing - simulates Azure Event Hub not returning it }, @@ -1044,7 +1044,7 @@ func TestGetLagForPartition_MissingPartition(t *testing.T) { { name: "Scenario 2: scaleToZeroOnInvalidOffset false, invalid consumer offset, missing partition", consumerOffset: invalidOffset, - topicPartitionOffsets: map[string]map[int32]int64{ + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ "test-topic": { // Partition 0 is missing }, @@ -1059,7 +1059,7 @@ func TestGetLagForPartition_MissingPartition(t *testing.T) { { name: "Scenario 3: Valid consumer offset, missing partition in topicPartitionOffsets", consumerOffset: 100, - topicPartitionOffsets: map[string]map[int32]int64{ + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ "test-topic": { // Partition 0 is missing - consumer has offset 100 but we can't get latestOffset }, @@ -1074,9 +1074,9 @@ func TestGetLagForPartition_MissingPartition(t *testing.T) { { name: "Control: Valid offsets, no missing partition", consumerOffset: 100, - topicPartitionOffsets: map[string]map[int32]int64{ + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ "test-topic": { - 0: 150, // Latest offset exists + 0: {latestOffset: 150}, // Latest offset exists }, }, offsetResetPolicy: earliest, @@ -1086,12 +1086,42 @@ func TestGetLagForPartition_MissingPartition(t *testing.T) { expectedError: false, description: "Normal case: both offsets exist, lag calculated correctly", }, + { + name: "Control: Invalid consumer offset with earliest policy uses retained log window", + consumerOffset: invalidOffset, + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ + "test-topic": { + 0: {earliestOffset: 100, earliestOffsetFound: true, latestOffset: 150}, + }, + }, + offsetResetPolicy: earliest, + scaleToZeroOnInvalidOffset: false, + expectedLag: 50, + expectedLagWithPersistent: 50, + expectedError: false, + description: "Invalid offset with earliest policy should return retained lag from log start to log end", + }, + { + name: "Control: Invalid consumer offset with earliest policy falls back to latest offset when earliest offset is unavailable", + consumerOffset: invalidOffset, + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ + "test-topic": { + 0: {latestOffset: 150}, + }, + }, + offsetResetPolicy: earliest, + scaleToZeroOnInvalidOffset: false, + expectedLag: 150, + expectedLagWithPersistent: 150, + expectedError: false, + description: "Invalid offset with earliest policy and no earliest offset should fall back to latest offset", + }, { name: "Control: Invalid consumer offset with latest policy and scaleToZeroOnInvalidOffset true", consumerOffset: invalidOffset, - topicPartitionOffsets: map[string]map[int32]int64{ + topicPartitionOffsets: map[string]map[int32]partitionOffsets{ "test-topic": { - 0: 150, + 0: {latestOffset: 150}, }, }, offsetResetPolicy: latest,