Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -203,6 +203,9 @@ manager: generate fmt vet ## Generate (kubebuilder) and build manager binary.
run: generate fmt vet
go run ./main.go

run-dev:
go run ./main.go --disable-cert-signing-support --disable-webhooks 2>&1 | jq -Cc -R 'fromjson? // empty'

# Install CRDs into a cluster by manually creating or replacing the CRD depending on whether is currently existing
# Apply is not applicable as the last-applied-configuration annotation would exceed the size limit enforced by the api server
install: manifests ## Install generated CRDs into the configured Kubernetes cluster.
Expand Down
4 changes: 2 additions & 2 deletions api/assets/kafka/jmx-exporter.yml
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,8 @@ rules:
# END: Export Kraft metrics

- pattern: 'kafka.server<type=(app-info), id=(\d+)><>(Version): ([-.~+\w\d]+)'
name: kafka_server_$1_$3
type: COUNTER
name: kafka_server_app_info_version
type: GAUGE
labels:
broker_id: $2
version: $4
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ spec:
selector:
matchLabels:
control-plane: controller-manager
component: alertmanager
app.kubernetes.io/component: operator
portLevelMtls:
{{ .Values.alertManager.port | quote }}:
mode: PERMISSIVE
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -133,8 +133,6 @@ spec:
app.kubernetes.io/name: {{ include "kafka-operator.name" . }}
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/component: operator
app: prometheus
component: alertmanager
spec:
{{- with .Values.imagePullSecrets }}
imagePullSecrets:
Expand Down
31 changes: 31 additions & 0 deletions charts/kafka-operator/templates/podmonitor.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
{{- if .Values.prometheusMetrics.podMonitor.enabled }}
kind: PodMonitor
apiVersion: monitoring.coreos.com/v1
metadata:
name: {{ include "kafka-operator.fullname" . }}
namespace: {{ .Release.Namespace | quote }}
labels:
helm.sh/chart: {{ include "kafka-operator.chart" . }}
app.kubernetes.io/name: {{ include "kafka-operator.name" . }}
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/managed-by: {{ .Release.Service }}
app.kubernetes.io/version: {{ .Chart.AppVersion }}
app.kubernetes.io/component: operator
{{- with .Values.operator.annotations }}
annotations:
{{- toYaml . | nindent 4 }}
{{- end }}
spec:
namespaceSelector:
matchNames:
- {{ .Release.Namespace }}
selector:
matchLabels:
app.kubernetes.io/name: {{ include "kafka-operator.name" . }}
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/component: operator
endpoints:
- interval: {{ .Values.prometheusMetrics.podMonitor.interval }}
port: metrics
path: /metrics
{{- end }}
4 changes: 4 additions & 0 deletions charts/kafka-operator/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,10 @@ prometheusMetrics:
create: true
# -- ServiceAccount used by prometheus auth proxy
name: kafka-operator-authproxy
podMonitor:
# -- If true, create a PodMonitor for Prometheus metrics
enabled: false
interval: 30s

# -- Health probes configuration
healthProbes: {}
Expand Down
42 changes: 31 additions & 11 deletions config/samples/kraft/simplekafkacluster_kraft.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,14 @@ metadata:
spec:
kRaft: true
monitoringConfig:
jmxImage: "ghcr.io/adobe/koperator/jmx-javaagent:1.4.0"
jmxImage: "ghcr.io/adobe/koperator/jmx-javaagent:1.5.0"
headlessServiceEnabled: true
propagateLabels: false
oneBrokerPerNode: false
clusterImage: "ghcr.io/adobe/koperator/kafka:2.13-3.9.1"
clusterImage: "ghcr.io/adobe/koperator/kafka:2.13-3.9.2-jdk21.0.11"
taintedBrokersSelector:
matchLabels:
restart: me
readOnlyConfig: |
auto.create.topics.enable=false
cruise.control.metrics.topic.auto.create=true
Expand All @@ -30,6 +33,9 @@ spec:
requests:
storage: 10Gi
broker:
kafkaJvmPerfOpts: >-
-Dcom.sun.management.jmxremote.port=1090 -Dcom.sun.management.jmxremote.rmi.port=1090
-Dcom.sun.management.jmxremote.local.only=false -Djava.rmi.server.hostname=127.0.0.1 -Djute.maxbuffer=0x9fffff
processRoles:
- broker
storageConfigs:
Expand All @@ -45,27 +51,26 @@ spec:
prometheus.io/port: "9020"
brokers:
- id: 0
brokerConfigGroup: "broker"
- id: 1
brokerConfigGroup: "broker"
- id: 2
brokerConfigGroup: "broker"
- id: 3
brokerConfigGroup: "default"
brokerConfig:
processRoles:
- controller
# - broker
- id: 4
- id: 1
brokerConfigGroup: "default"
brokerConfig:
processRoles:
- controller
- id: 5
- id: 2
brokerConfigGroup: "default"
brokerConfig:
processRoles:
- controller
- id: 100
brokerConfigGroup: "broker"
- id: 101
brokerConfigGroup: "broker"
- id: 102
brokerConfigGroup: "broker"
rollingUpgradeConfig:
failureThreshold: 1
listenersConfig:
Expand Down Expand Up @@ -300,3 +305,18 @@ spec:
{
"min.insync.replicas": 3
}
capacityConfig: |-
{
"brokerCapacities":[
{
"brokerId": "-1",
"capacity": {
"DISK": {"/kafka-logs-broker/kafka": "10240"},
"CPU": {"num.cores": "1"},
"NW_IN": "900000",
"NW_OUT": "900000"
},
"doc": "This is the default capacity. Capacity unit used for disk is in MB, cpu is in cores, network throughput is in KB."
}
]
}
38 changes: 37 additions & 1 deletion controllers/kafkacluster_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ package controllers

import (
"context"
"encoding/json"
"fmt"
"reflect"
"time"
Expand Down Expand Up @@ -359,6 +360,41 @@ func (r *KafkaClusterReconciler) updateAndFetchLatest(ctx context.Context, clust
return cluster, nil
}

// ignoreMetadataFields returns a patch.CalculateOption that drops only the given fields
// under metadata (e.g. "managedFields", "resourceVersion") instead of the whole metadata
// object like patch.IgnoreField("metadata") does, so other metadata changes (labels,
// annotations, etc.) still count toward the diff.
func ignoreMetadataFields(fields ...string) patch.CalculateOption {
return func(current, modified []byte) ([]byte, []byte, error) {
current, err := deleteMetadataFields(current, fields...)
if err != nil {
return nil, nil, errors.WrapIf(err, "could not delete metadata fields from current byte sequence")
}

modified, err = deleteMetadataFields(modified, fields...)
if err != nil {
return nil, nil, errors.WrapIf(err, "could not delete metadata fields from modified byte sequence")
}

return current, modified, nil
}
}

func deleteMetadataFields(obj []byte, fields ...string) ([]byte, error) {
var objectMap map[string]interface{}
if err := json.Unmarshal(obj, &objectMap); err != nil {
return nil, errors.WrapIf(err, "could not unmarshal byte sequence")
}

if metadata, ok := objectMap["metadata"].(map[string]interface{}); ok {
for _, field := range fields {
delete(metadata, field)
}
}

return json.Marshal(objectMap)
}

// SetupKafkaClusterWithManager registers kafka cluster controller to the manager
func SetupKafkaClusterWithManager(mgr ctrl.Manager, contourEnabled bool) *ctrl.Builder {
log := mgr.GetLogger()
Expand Down Expand Up @@ -390,7 +426,7 @@ func SetupKafkaClusterWithManager(mgr ctrl.Manager, contourEnabled bool) *ctrl.B
UpdateFunc: func(e event.UpdateEvent) bool {
switch newObj := e.ObjectNew.(type) {
case *corev1.Pod, *corev1.ConfigMap, *corev1.PersistentVolumeClaim:
patchResult, err := patch.DefaultPatchMaker.Calculate(e.ObjectOld, e.ObjectNew)
patchResult, err := patch.DefaultPatchMaker.Calculate(e.ObjectOld, e.ObjectNew, patch.IgnoreStatusFields(), ignoreMetadataFields("managedFields", "resourceVersion"))
if err != nil {
log.Error(err, "could not match objects", "kind", e.ObjectOld.GetObjectKind())
} else if patchResult.IsEmpty() {
Expand Down
10 changes: 6 additions & 4 deletions pkg/jmxextractor/extractor.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"io"
"net/http"
"regexp"
"time"

"github.com/banzaicloud/koperator/api/v1beta1"
"github.com/banzaicloud/koperator/pkg/errorfactory"
Expand All @@ -30,9 +31,9 @@ import (
)

const (
headlessServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessServiceTemplate + ".%s.svc.%s:%d"
headlessControllerServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessControllerServiceTemplate + ".%s.svc.%s:%d"
serviceJMXTemplate = "http://%s-%d.%s.svc.%s:%d"
headlessServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessServiceTemplate + ".%s.svc.%s:%d/metrics"
headlessControllerServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessControllerServiceTemplate + ".%s.svc.%s:%d/metrics"
serviceJMXTemplate = "http://%s-%d.%s.svc.%s:%d/metrics"
versionRegexGroup = "version"
)

Expand Down Expand Up @@ -90,7 +91,8 @@ func (exp *jmxExtractor) ExtractDockerImageAndVersion(brokerId int32, brokerConf
requestURL = fmt.Sprintf(serviceJMXTemplate, exp.clusterName, brokerId,
exp.clusterNamespace, exp.kubernetesClusterDomain, 9020)
}
rsp, err := http.Get(requestURL)
client := &http.Client{Timeout: 30 * time.Second}
rsp, err := client.Get(requestURL)
if err != nil {
exp.log.Error(err, fmt.Sprintf("error during talking to broker-%d", brokerId))
return nil, errorfactory.New(errorfactory.BrokersNotReady{}, err, "unable to talk to ...")
Expand Down
2 changes: 1 addition & 1 deletion pkg/k8sutil/status.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ func UpdateBrokerStatus(c client.Client, brokerIDs []string, cluster *banzaiclou
}
// update loses the typeMeta of the config that's used later when setting ownerrefs
cluster.TypeMeta = typeMeta
logger.Info("Kafka cluster state updated")
logger.V(1).Info("Broker status updated", "brokers", brokerIDs, "status", state)
return nil
}

Expand Down
2 changes: 1 addition & 1 deletion pkg/kafkaclient/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -130,7 +130,7 @@ func NewFromCluster(k8sclient client.Client, cluster *v1beta1.KafkaCluster) (Kaf
if err := client.Close(); err != nil {
log.Error(err, "Error closing Kafka client")
} else {
log.Info("Kafka client closed cleanly")
log.V(1).Info("Kafka client closed cleanly")
}
}
return client, close, err
Expand Down
6 changes: 3 additions & 3 deletions pkg/resources/cruisecontrol/configmap.go
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,7 @@ type JBODInvariantCapacityConfig struct {
func GenerateCapacityConfig(kafkaCluster *v1beta1.KafkaCluster, log logr.Logger, config *corev1.ConfigMap) (string, error) {
var err error

log.Info("generating capacity config")
log.V(2).Info("generating capacity config")

var capacityConfig JBODInvariantCapacityConfig
var userConfigBrokerIds []string
Expand All @@ -184,7 +184,7 @@ func GenerateCapacityConfig(kafkaCluster *v1beta1.KafkaCluster, log logr.Logger,
}
// If the -1 default exists we don't have to do anything else here since all brokers will have values.
if brokerId == "-1" {
log.Info("Using user provided capacity config because it has universal default defined", "capacity config", userProvidedCapacityConfig)
log.V(2).Info("Using user provided capacity config because it has universal default defined", "capacity config", userProvidedCapacityConfig)
return userProvidedCapacityConfig, nil
}
userConfigBrokerIds = append(userConfigBrokerIds, brokerId)
Expand All @@ -211,7 +211,7 @@ func GenerateCapacityConfig(kafkaCluster *v1beta1.KafkaCluster, log logr.Logger,
if err != nil {
return "", errors.WrapIf(err, "could not marshal cruise control capacity config")
}
log.Info("broker capacity config generated successfully", "capacity config", string(result))
log.V(2).Info("broker capacity config generated successfully", "capacity config", string(result))
return string(result), err
}

Expand Down
Loading