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
7 changes: 7 additions & 0 deletions api/assets/kafka/kraft-controller-healthcheck.sh
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,13 @@ fi
# Reachable and reporting some other raft state (e.g. candidate, unattached, observer).
if echo "${METRICS}" | grep -q "^${METRIC_PREFIX}"; then
STATE=$(echo "${METRICS}" | grep "^${METRIC_PREFIX}" | head -n 1 | sed -E "s/^${METRIC_PREFIX}([a-z]+).*/\1/")
# Under a dynamic controller quorum, a controller that has not yet been promoted with
# "kafka-metadata-quorum add-controller" will be a non-voting observer until promoted.
if [ "${STATE}" = "observer" ]; then
echo "The controller is a non-voting observer (dynamic quorum: awaiting add-controller)."
[ "${mode}" = "readiness" ] && exit 1
exit 0
fi
echo "Failure: the controller is in an unexpected state: ${STATE}. Expecting 'leader' or 'follower'."
exit 1
fi
Expand Down
5 changes: 5 additions & 0 deletions api/v1beta1/kafkacluster_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,11 @@ type KafkaClusterStatus struct {
ListenerStatuses ListenerStatuses `json:"listenerStatuses,omitempty"`
// ClusterID is a base64-encoded random UUID generated by Koperator to run the Kafka cluster in KRaft mode
ClusterID string `json:"clusterID,omitempty"`
// KRaftDynamicQuorumBootstrapped is set to true once the dynamic KRaft controller quorum has formed
// (a controller has reported a leader/follower raft state). Once set, Koperator stops formatting the
// lowest-ID controller with --standalone, so a controller that later loses its disk rejoins the
// existing quorum as an observer instead of bootstrapping a divergent single-voter quorum.
KRaftDynamicQuorumBootstrapped bool `json:"kRaftDynamicQuorumBootstrapped,omitempty"`
}

// RollingUpgradeStatus defines status of rolling upgrade
Expand Down
7 changes: 7 additions & 0 deletions charts/kafka-operator/crds/kafkaclusters.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -23845,6 +23845,13 @@ spec:
description: CruiseControlTopicStatus holds info about the CC topic
status
type: string
kRaftDynamicQuorumBootstrapped:
description: |-
KRaftDynamicQuorumBootstrapped is set to true once the dynamic KRaft controller quorum has formed
(a controller has reported a leader/follower raft state). Once set, Koperator stops formatting the
lowest-ID controller with --standalone, so a controller that later loses its disk rejoins the
existing quorum as an observer instead of bootstrapping a divergent single-voter quorum.
type: boolean
listenerStatuses:
description: |-
ListenerStatuses holds information about the statuses of the configured listeners.
Expand Down
7 changes: 7 additions & 0 deletions config/base/crds/kafka.banzaicloud.io_kafkaclusters.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -23845,6 +23845,13 @@ spec:
description: CruiseControlTopicStatus holds info about the CC topic
status
type: string
kRaftDynamicQuorumBootstrapped:
description: |-
KRaftDynamicQuorumBootstrapped is set to true once the dynamic KRaft controller quorum has formed
(a controller has reported a leader/follower raft state). Once set, Koperator stops formatting the
lowest-ID controller with --standalone, so a controller that later loses its disk rejoins the
existing quorum as an observer instead of bootstrapping a divergent single-voter quorum.
type: boolean
listenerStatuses:
description: |-
ListenerStatuses holds information about the statuses of the configured listeners.
Expand Down
290 changes: 290 additions & 0 deletions config/samples/kraft/simplekafkacluster_kraft_dynamic.yaml

Large diffs are not rendered by default.

33 changes: 24 additions & 9 deletions pkg/resources/kafka/configmap.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ import (
properties "github.com/banzaicloud/koperator/properties/pkg"
)

func (r *Reconciler) getConfigProperties(bConfig *v1beta1.BrokerConfig, broker v1beta1.Broker, quorumVoters []string,
func (r *Reconciler) getConfigProperties(bConfig *v1beta1.BrokerConfig, broker v1beta1.Broker, quorumVoters, quorumBootstrapServers []string,
extListenerStatuses, intListenerStatuses, controllerIntListenerStatuses map[string]v1beta1.ListenerStatusList,
serverPasses map[string]string, clientPass string, superUsers []string, log logr.Logger) *properties.Properties {
config := properties.NewProperties()
Expand Down Expand Up @@ -69,7 +69,7 @@ func (r *Reconciler) getConfigProperties(bConfig *v1beta1.BrokerConfig, broker v

// Kafka Broker configurations
if r.KafkaCluster.Spec.KRaftMode {
configureBrokerKRaftMode(bConfig, broker.Id, r.KafkaCluster, config, quorumVoters, serverPasses, extListenerStatuses, intListenerStatuses, log, brokerReadOnlyConfig)
configureBrokerKRaftMode(bConfig, broker.Id, r.KafkaCluster, config, quorumVoters, quorumBootstrapServers, serverPasses, extListenerStatuses, intListenerStatuses, log, brokerReadOnlyConfig)
} else {
configureBrokerZKMode(broker.Id, r.KafkaCluster, config, extListenerStatuses, intListenerStatuses, controllerIntListenerStatuses, log)
}
Expand Down Expand Up @@ -157,7 +157,7 @@ func (r *Reconciler) configCCMetricsReporter(broker v1beta1.Broker, bConfig *v1b
}

func configureBrokerKRaftMode(bConfig *v1beta1.BrokerConfig, brokerID int32, kafkaCluster *v1beta1.KafkaCluster, config *properties.Properties,
quorumVoters []string, serverPasses map[string]string, extListenerStatuses, intListenerStatuses map[string]v1beta1.ListenerStatusList, log logr.Logger,
quorumVoters, quorumBootstrapServers []string, serverPasses map[string]string, extListenerStatuses, intListenerStatuses map[string]v1beta1.ListenerStatusList, log logr.Logger,
brokerReadOnlyConfig *properties.Properties) {
controllerListenerName := generateControlPlaneListener(kafkaCluster.Spec.ListenersConfig.InternalListeners)

Expand Down Expand Up @@ -185,8 +185,14 @@ func configureBrokerKRaftMode(bConfig *v1beta1.BrokerConfig, brokerID int32, kaf
}

if shouldConfigureControllerQuorumForBroker(brokerReadOnlyConfig) {
if err := config.Set(kafkautils.KafkaConfigControllerQuorumVoters, quorumVoters); err != nil {
log.Error(err, fmt.Sprintf(kafkautils.BrokerConfigErrorMsgTemplate, kafkautils.KafkaConfigControllerQuorumVoters))
if shouldUseDynamicKRaftQuorum(brokerReadOnlyConfig) {
if err := config.Set(kafkautils.KafkaConfigControllerQuorumBootstrapServers, quorumBootstrapServers); err != nil {
log.Error(err, fmt.Sprintf(kafkautils.BrokerConfigErrorMsgTemplate, kafkautils.KafkaConfigControllerQuorumBootstrapServers))
}
} else {
if err := config.Set(kafkautils.KafkaConfigControllerQuorumVoters, quorumVoters); err != nil {
log.Error(err, fmt.Sprintf(kafkautils.BrokerConfigErrorMsgTemplate, kafkautils.KafkaConfigControllerQuorumVoters))
}
}

if controllerListenerName != "" {
Expand Down Expand Up @@ -255,6 +261,14 @@ func shouldConfigureControllerQuorumForBroker(brokerReadOnlyConfig *properties.P
return !found || migrationBrokerControllerQuorumConfigEnabled.Value() == configValueTrue
}

// shouldUseDynamicKRaftQuorum returns true only when DynamicKRaftControllerQuorum is set
// to 'true'. It defaults to false (static quorum) when absent, so
// existing static-quorum clusters will continue to remain static.
func shouldUseDynamicKRaftQuorum(brokerReadOnlyConfig *properties.Properties) bool {
dynamicKRaftControllerQuorum, found := brokerReadOnlyConfig.Get(kafkautils.DynamicKRaftControllerQuorum)
return found && dynamicKRaftControllerQuorum.Value() == configValueTrue
}

func configureBrokerZKMode(brokerID int32, kafkaCluster *v1beta1.KafkaCluster, config *properties.Properties, extListenerStatuses, intListenerStatuses,
controllerIntListenerStatuses map[string]v1beta1.ListenerStatusList, log logr.Logger) {
if err := config.Set(kafkautils.KafkaConfigBrokerID, brokerID); err != nil {
Expand Down Expand Up @@ -351,7 +365,7 @@ func generateSuperUsers(users []string) (suStrings []string) {
return
}

func (r *Reconciler) configMap(broker v1beta1.Broker, brokerConfig *v1beta1.BrokerConfig, quorumVoters []string,
func (r *Reconciler) configMap(broker v1beta1.Broker, brokerConfig *v1beta1.BrokerConfig, quorumVoters, quorumBootstrapServers []string,
extListenerStatuses, intListenerStatuses, controllerIntListenerStatuses map[string]v1beta1.ListenerStatusList,
serverPasses map[string]string, clientPass string, superUsers []string, log logr.Logger) *corev1.ConfigMap {
brokerConf := &corev1.ConfigMap{
Expand All @@ -363,7 +377,7 @@ func (r *Reconciler) configMap(broker v1beta1.Broker, brokerConfig *v1beta1.Brok
),
r.KafkaCluster,
),
Data: map[string]string{kafkautils.ConfigPropertyName: r.generateBrokerConfig(broker, brokerConfig, quorumVoters, extListenerStatuses,
Data: map[string]string{kafkautils.ConfigPropertyName: r.generateBrokerConfig(broker, brokerConfig, quorumVoters, quorumBootstrapServers, extListenerStatuses,
intListenerStatuses, controllerIntListenerStatuses, serverPasses, clientPass, superUsers, log)},
}
if brokerConfig.Log4jConfig != "" {
Expand Down Expand Up @@ -612,13 +626,13 @@ func mergeSuperUsersPropertyValue(source *properties.Properties, target *propert
return ""
}

func (r Reconciler) generateBrokerConfig(broker v1beta1.Broker, brokerConfig *v1beta1.BrokerConfig, quorumVoters []string,
func (r Reconciler) generateBrokerConfig(broker v1beta1.Broker, brokerConfig *v1beta1.BrokerConfig, quorumVoters, quorumBootstrapServers []string,
extListenerStatuses, intListenerStatuses, controllerIntListenerStatuses map[string]v1beta1.ListenerStatusList,
serverPasses map[string]string, clientPass string, superUsers []string, log logr.Logger) string {
finalBrokerConfig := getBrokerReadOnlyConfig(broker, r.KafkaCluster, log)

// Get operator generated configuration
opGenConf := r.getConfigProperties(brokerConfig, broker, quorumVoters, extListenerStatuses, intListenerStatuses,
opGenConf := r.getConfigProperties(brokerConfig, broker, quorumVoters, quorumBootstrapServers, extListenerStatuses, intListenerStatuses,
controllerIntListenerStatuses, serverPasses, clientPass, superUsers, log)

// Merge operator generated configuration to the final one
Expand All @@ -636,6 +650,7 @@ func (r Reconciler) generateBrokerConfig(broker v1beta1.Broker, brokerConfig *v1
// Remove the migration broker configuration since its only used as flags to derive other configs
finalBrokerConfig.Delete(kafkautils.MigrationBrokerControllerQuorumConfigEnabled)
finalBrokerConfig.Delete(kafkautils.MigrationBrokerKRaftMode)
finalBrokerConfig.Delete(kafkautils.DynamicKRaftControllerQuorum)

finalBrokerConfig.Sort()

Expand Down
168 changes: 165 additions & 3 deletions pkg/resources/kafka/configmap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -802,7 +802,7 @@ zookeeper.connect=example.zk:2181/`,
superUsers = []string{"CN=kafka-headless.kafka.svc.cluster.local"}
}

generatedConfig := r.generateBrokerConfig(r.KafkaCluster.Spec.Brokers[0], r.KafkaCluster.Spec.Brokers[0].BrokerConfig, nil, map[string]v1beta1.ListenerStatusList{},
generatedConfig := r.generateBrokerConfig(r.KafkaCluster.Spec.Brokers[0], r.KafkaCluster.Spec.Brokers[0].BrokerConfig, nil, nil, map[string]v1beta1.ListenerStatusList{},
map[string]v1beta1.ListenerStatusList{}, controllerListenerStatus, serverPasses, clientPass, superUsers, logr.Discard())

generated, err := properties.NewFromString(generatedConfig)
Expand Down Expand Up @@ -1296,6 +1296,160 @@ log.dirs=/test-kafka-logs/kafka
metric.reporters=com.linkedin.kafka.cruisecontrol.metricsreporter.CruiseControlMetricsReporter
node.id=300
process.roles=broker
`},
},
{
testName: "a Kafka cluster with dynamic KRaft controller quorum enabled renders controller.quorum.bootstrap.servers instead of controller.quorum.voters",
brokers: []v1beta1.Broker{
{
Id: 0,
BrokerConfig: &v1beta1.BrokerConfig{
Roles: []string{"broker"},
StorageConfigs: []v1beta1.StorageConfig{{MountPath: "/test-kafka-logs"}},
},
ReadOnlyConfig: "kraft.dynamicControllerQuorum.enabled=true",
},
{
Id: 50,
BrokerConfig: &v1beta1.BrokerConfig{
Roles: []string{"controller"},
StorageConfigs: []v1beta1.StorageConfig{{MountPath: "/test-kafka-logs"}},
},
ReadOnlyConfig: "kraft.dynamicControllerQuorum.enabled=true",
},
},
listenersConfig: v1beta1.ListenersConfig{
InternalListeners: []v1beta1.InternalListenerConfig{
{
CommonListenerSpec: v1beta1.CommonListenerSpec{
Type: v1beta1.SecurityProtocol("PLAINTEXT"),
Name: "internal",
ContainerPort: 9092,
UsedForInnerBrokerCommunication: true,
},
},
{
CommonListenerSpec: v1beta1.CommonListenerSpec{
Type: v1beta1.SecurityProtocol("PLAINTEXT"),
Name: "controller",
ContainerPort: 9093,
},
UsedForControllerCommunication: true,
},
},
},
internalListenerStatuses: map[string]v1beta1.ListenerStatusList{
"internal": {
{Name: "broker-0", Address: "kafka-0.kafka.svc.cluster.local:9092"},
{Name: "broker-50", Address: "kafka-50.kafka.svc.cluster.local:9092"},
},
},
controllerListenerStatus: map[string]v1beta1.ListenerStatusList{
"controller": {
{Name: "broker-0", Address: "kafka-0.kafka.svc.cluster.local:9093"},
{Name: "broker-50", Address: "kafka-50.kafka.svc.cluster.local:9093"},
},
},
expectedBrokerConfigs: []string{
`advertised.listeners=INTERNAL://kafka-0.kafka.svc.cluster.local:9092
controller.listener.names=CONTROLLER
controller.quorum.bootstrap.servers=kafka-50.kafka.svc.cluster.local:9093
cruise.control.metrics.reporter.bootstrap.servers=kafka-all-broker.kafka.svc.cluster.local:9092
cruise.control.metrics.reporter.kubernetes.mode=true
inter.broker.listener.name=INTERNAL
listener.security.protocol.map=INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT
listeners=INTERNAL://:9092
log.dirs=/test-kafka-logs/kafka
metric.reporters=com.linkedin.kafka.cruisecontrol.metricsreporter.CruiseControlMetricsReporter
node.id=0
process.roles=broker
`,
`controller.listener.names=CONTROLLER
controller.quorum.bootstrap.servers=kafka-50.kafka.svc.cluster.local:9093
inter.broker.listener.name=INTERNAL
listener.security.protocol.map=INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT
listeners=CONTROLLER://:9093
log.dirs=/test-kafka-logs/kafka
node.id=50
process.roles=controller
`},
},
{
testName: "dynamic KRaft controller quorum flag composes with the ZK-migration bridge flags (broker still bridging to ZK, controller fully on dynamic quorum)",
brokers: []v1beta1.Broker{
{
Id: 0,
BrokerConfig: &v1beta1.BrokerConfig{
Roles: []string{"broker"},
StorageConfigs: []v1beta1.StorageConfig{{MountPath: "/test-kafka-logs"}},
},
ReadOnlyConfig: "kraft.dynamicControllerQuorum.enabled=true\nmigration.broker.kRaftMode=false",
},
{
Id: 50,
BrokerConfig: &v1beta1.BrokerConfig{
Roles: []string{"controller"},
StorageConfigs: []v1beta1.StorageConfig{{MountPath: "/test-kafka-logs"}},
},
ReadOnlyConfig: "kraft.dynamicControllerQuorum.enabled=true",
},
},
listenersConfig: v1beta1.ListenersConfig{
InternalListeners: []v1beta1.InternalListenerConfig{
{
CommonListenerSpec: v1beta1.CommonListenerSpec{
Type: v1beta1.SecurityProtocol("PLAINTEXT"),
Name: "internal",
ContainerPort: 9092,
UsedForInnerBrokerCommunication: true,
},
},
{
CommonListenerSpec: v1beta1.CommonListenerSpec{
Type: v1beta1.SecurityProtocol("PLAINTEXT"),
Name: "controller",
ContainerPort: 9093,
},
UsedForControllerCommunication: true,
},
},
},
internalListenerStatuses: map[string]v1beta1.ListenerStatusList{
"internal": {
{Name: "broker-0", Address: "kafka-0.kafka.svc.cluster.local:9092"},
{Name: "broker-50", Address: "kafka-50.kafka.svc.cluster.local:9092"},
},
},
controllerListenerStatus: map[string]v1beta1.ListenerStatusList{
"controller": {
{Name: "broker-0", Address: "kafka-0.kafka.svc.cluster.local:9093"},
{Name: "broker-50", Address: "kafka-50.kafka.svc.cluster.local:9093"},
},
},
zkAddresses: []string{"example.zk:2181"},
zkPath: "/kafka",
expectedBrokerConfigs: []string{
`advertised.listeners=INTERNAL://kafka-0.kafka.svc.cluster.local:9092
broker.id=0
controller.listener.names=CONTROLLER
controller.quorum.bootstrap.servers=kafka-50.kafka.svc.cluster.local:9093
cruise.control.metrics.reporter.bootstrap.servers=kafka-all-broker.kafka.svc.cluster.local:9092
cruise.control.metrics.reporter.kubernetes.mode=true
inter.broker.listener.name=INTERNAL
listener.security.protocol.map=INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT
listeners=INTERNAL://:9092
log.dirs=/test-kafka-logs/kafka
metric.reporters=com.linkedin.kafka.cruisecontrol.metricsreporter.CruiseControlMetricsReporter
zookeeper.connect=example.zk:2181/kafka
`,
`controller.listener.names=CONTROLLER
controller.quorum.bootstrap.servers=kafka-50.kafka.svc.cluster.local:9093
inter.broker.listener.name=INTERNAL
listener.security.protocol.map=INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT
listeners=CONTROLLER://:9093
log.dirs=/test-kafka-logs/kafka
node.id=50
process.roles=controller
`},
},
}
Expand Down Expand Up @@ -1332,8 +1486,12 @@ process.roles=broker
if err != nil {
t.Error(err)
}
quorumBootstrapServers, err := generateQuorumBootstrapServers(r.KafkaCluster, test.controllerListenerStatus)
if err != nil {
t.Error(err)
}

generatedConfig := r.generateBrokerConfig(b, b.BrokerConfig, quorumVoters, map[string]v1beta1.ListenerStatusList{},
generatedConfig := r.generateBrokerConfig(b, b.BrokerConfig, quorumVoters, quorumBootstrapServers, map[string]v1beta1.ListenerStatusList{},
test.internalListenerStatuses, test.controllerListenerStatus, nil, "", nil, logr.Discard())

require.Equal(t, test.expectedBrokerConfigs[i], generatedConfig)
Expand Down Expand Up @@ -1647,7 +1805,11 @@ process.roles=controller
if err != nil {
t.Error(err)
}
generatedConfig := r.generateBrokerConfig(b, b.BrokerConfig, quorumVoters, map[string]v1beta1.ListenerStatusList{},
quorumBootstrapServers, err := generateQuorumBootstrapServers(r.KafkaCluster, test.controllerListenerStatus)
if err != nil {
t.Error(err)
}
generatedConfig := r.generateBrokerConfig(b, b.BrokerConfig, quorumVoters, quorumBootstrapServers, map[string]v1beta1.ListenerStatusList{},
test.internalListenerStatuses, test.controllerListenerStatus, nil, "", nil, logr.Discard())

require.Equal(t, test.expectedBrokerConfigs[i], generatedConfig)
Expand Down
Loading
Loading