Skip to content
Draft
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 README.md
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,9 @@ kubectl -n kafka run kafka-consumer -it --image=adobe/kafka:2.13-3.9.1 --rm=true

## Documentation

For opt-in, broker-only KRaft metadata PVCs with an automatic, one-broker-at-a-time live migration,
see [dedicated broker metadata storage](docs/broker-metadata-storage.md).

For detailed documentation on the Koperator project, see the [Koperator documentation website](https://opensource.adobe.com/koperator/).

## Issues and contributions
Expand Down
12 changes: 12 additions & 0 deletions api/v1beta1/common_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,9 @@ type SSLClientAuthentication string
// PerBrokerConfigurationState holds info about the per-broker configuration state
type PerBrokerConfigurationState string

// MetadataStorageState holds info about the dedicated KRaft metadata storage of a broker
type MetadataStorageState string

// ExternalListenerConfigNames type describes a collection of external listener names
type ExternalListenerConfigNames []string

Expand Down Expand Up @@ -253,8 +256,17 @@ type BrokerState struct {
Image string `json:"image,omitempty"`
// Compressed data from broker configuration to restore broker pod in specific cases
ConfigurationBackup string `json:"configurationBackup,omitempty"`
// MetadataStorageState is Ready once the broker runs from its dedicated metadata storage.
// Data disks of a broker with metadataStorage can be removed only after that.
MetadataStorageState MetadataStorageState `json:"metadataStorageState,omitempty"`
}

const (
// MetadataStorageReady states that a replacement broker pod became ready using its
// dedicated metadata storage, so the metadata no longer lives on any data disk
MetadataStorageReady MetadataStorageState = "Ready"
)

const (
// Configured states the broker is running
Configured RackAwarenessState = "Configured"
Expand Down
41 changes: 29 additions & 12 deletions api/v1beta1/kafkacluster_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -301,18 +301,23 @@ type BrokerConfig struct {
// +kubebuilder:validation:Items:Type=string
// +kubebuilder:validation:Items:Enum=controller;broker
// +optional
Roles []string `json:"processRoles,omitempty"`
Image string `json:"image,omitempty"`
MetricsReporterImage string `json:"metricsReporterImage,omitempty"`
Config string `json:"config,omitempty"`
StorageConfigs []StorageConfig `json:"storageConfigs,omitempty"`
ServiceAccountName string `json:"serviceAccountName,omitempty"`
Resources *corev1.ResourceRequirements `json:"resourceRequirements,omitempty"`
ImagePullSecrets []corev1.LocalObjectReference `json:"imagePullSecrets,omitempty"`
NodeSelector map[string]string `json:"nodeSelector,omitempty"`
Tolerations []corev1.Toleration `json:"tolerations,omitempty"`
KafkaHeapOpts string `json:"kafkaHeapOpts,omitempty"`
KafkaJVMPerfOpts string `json:"kafkaJvmPerfOpts,omitempty"`
Roles []string `json:"processRoles,omitempty"`
Image string `json:"image,omitempty"`
MetricsReporterImage string `json:"metricsReporterImage,omitempty"`
Config string `json:"config,omitempty"`
StorageConfigs []StorageConfig `json:"storageConfigs,omitempty"`
// MetadataStorage optionally isolates KRaft metadata on a dedicated PVC.
// Only broker-only KRaft nodes may use it. It is not a data or Cruise Control disk.
// Once enabled it cannot be removed or relocated; reverse migration is not supported.
// +optional
MetadataStorage *StorageConfig `json:"metadataStorage,omitempty"`
ServiceAccountName string `json:"serviceAccountName,omitempty"`
Resources *corev1.ResourceRequirements `json:"resourceRequirements,omitempty"`
ImagePullSecrets []corev1.LocalObjectReference `json:"imagePullSecrets,omitempty"`
NodeSelector map[string]string `json:"nodeSelector,omitempty"`
Tolerations []corev1.Toleration `json:"tolerations,omitempty"`
KafkaHeapOpts string `json:"kafkaHeapOpts,omitempty"`
KafkaJVMPerfOpts string `json:"kafkaJvmPerfOpts,omitempty"`
// Override for the default log4j configuration
Log4jConfig string `json:"log4jConfig,omitempty"`
// Custom annotations for the broker pods - e.g.: Prometheus scraping annotations:
Expand Down Expand Up @@ -1300,12 +1305,24 @@ func (b *Broker) GetBrokerConfig(kafkaClusterSpec KafkaClusterSpec) (*BrokerConf
}
envs := mergeEnvs(kafkaClusterSpec, &groupConfig, bConfig)

// A broker override replaces this single storage object, rather than merging
// half of a group PVC spec into a broker-local mount.
metadataStorage := bConfig.MetadataStorage
if metadataStorage == nil {
metadataStorage = groupConfig.MetadataStorage
}
if metadataStorage != nil {
metadataStorage = metadataStorage.DeepCopy()
}
err = mergo.Merge(bConfig, groupConfig, mergo.WithAppendSlice)
if err != nil {
return nil, errors.WrapIf(err, "could not merge brokerConfig with ConfigGroup")
}

bConfig.StorageConfigs = dedupStorageConfigs(bConfig.StorageConfigs)
if metadataStorage != nil {
bConfig.MetadataStorage = metadataStorage.DeepCopy()
}
if groupConfig.Affinity != nil || bConfig.Affinity != nil {
bConfig.Affinity = dstAffinity
}
Expand Down
92 changes: 92 additions & 0 deletions api/v1beta1/metadata_storage.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
// Copyright 2026 Adobe. All rights reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

package v1beta1

import (
"fmt"
"path"
"reflect"
"regexp"
"strings"
)

// MetadataStorageVolumeName is the pod volume name reserved for the dedicated metadata PVC.
const MetadataStorageVolumeName = "kraft-metadata"

var metadataMountPathPattern = regexp.MustCompile(`^/[a-zA-Z0-9_./-]+$`)

// ValidateMetadataStorage validates the effective (group-merged) broker configuration.
func (b *BrokerConfig) ValidateMetadataStorage(kraft bool) error {
if b == nil || b.MetadataStorage == nil {
return nil
}
if !kraft || !b.IsBrokerOnlyNode() || len(b.Roles) != 1 {
return fmt.Errorf("metadataStorage requires a broker-only KRaft node")
}
m := b.MetadataStorage
if m.PvcSpec == nil || m.EmptyDir != nil {
return fmt.Errorf("metadataStorage requires pvcSpec and does not support emptyDir")
}
if m.PvcSpec.VolumeMode != nil && *m.PvcSpec.VolumeMode != "Filesystem" {
return fmt.Errorf("metadataStorage requires a filesystem PVC")
}
if !metadataMountPathPattern.MatchString(m.MountPath) || path.Clean(m.MountPath) != m.MountPath || m.MountPath == "/" {
return fmt.Errorf("metadataStorage mountPath must be a clean absolute path using letters, digits, '/', '.', '_' or '-'")
}
paths := []string{"/config", "/opt", "/etc", "/var/run", "/run", "/dev", "/proc", "/sys", "/usr", "/bin", "/sbin", "/lib", "/lib64", "/tmp"}
for _, s := range b.StorageConfigs {
// Ephemeral data cannot be a migration source, and durable metadata gains nothing
// beside ephemeral replicas.
if s.PvcSpec == nil || s.EmptyDir != nil {
return fmt.Errorf("metadataStorage requires all data storageConfigs to be PVC-backed; %q is not", s.MountPath)
}
paths = append(paths, s.MountPath)
}
for _, v := range b.VolumeMounts {
if v.Name == MetadataStorageVolumeName {
return fmt.Errorf("%s is a reserved volume name", MetadataStorageVolumeName)
}
paths = append(paths, v.MountPath)
}
for _, v := range b.Volumes {
if v.Name == MetadataStorageVolumeName {
return fmt.Errorf("%s is a reserved volume name", MetadataStorageVolumeName)
}
}
for _, p := range paths {
if StoragePathsOverlap(m.MountPath, p) {
return fmt.Errorf("metadataStorage mountPath overlaps %q", p)
}
}
return nil
}

// StoragePathsOverlap includes nested mounts, which can hide an existing volume.
func StoragePathsOverlap(a, b string) bool {
a, b = path.Clean(a), path.Clean(b)
return a == b || strings.HasPrefix(a, b+"/") || strings.HasPrefix(b, a+"/")
}

// MetadataStorageLocationEqual allows capacity changes, but no relocation or PVC replacement.
func MetadataStorageLocationEqual(a, b *StorageConfig) bool {
if a == nil || b == nil {
return a == b
}
a, b = a.DeepCopy(), b.DeepCopy()
if a.PvcSpec != nil && b.PvcSpec != nil {
a.PvcSpec.Resources = b.PvcSpec.Resources
}
return reflect.DeepEqual(a, b)
}
98 changes: 98 additions & 0 deletions api/v1beta1/metadata_storage_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
// Copyright 2026 Adobe. All rights reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

package v1beta1

import (
"testing"

corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
)

func TestMetadataStorageValidationAndMapping(t *testing.T) {
base := &BrokerConfig{
Roles: []string{"broker"},
StorageConfigs: []StorageConfig{{MountPath: "/csi-kafka-logs1", PvcSpec: &corev1.PersistentVolumeClaimSpec{}}},
MetadataStorage: &StorageConfig{MountPath: "/csi-kafka-metadata", PvcSpec: &corev1.PersistentVolumeClaimSpec{}},
}
for _, tc := range []struct {
name string
kraft bool
change func(*BrokerConfig)
valid bool
}{
{"broker", true, func(*BrokerConfig) {}, true},
{"default unchanged", false, func(b *BrokerConfig) { b.MetadataStorage = nil }, true},
{"zk", false, func(*BrokerConfig) {}, false},
{"controller", true, func(b *BrokerConfig) { b.Roles = []string{"controller"} }, false},
{"combined", true, func(b *BrokerConfig) { b.Roles = []string{"broker", "controller"} }, false},
{"emptydir", true, func(b *BrokerConfig) { b.MetadataStorage.EmptyDir = &corev1.EmptyDirVolumeSource{} }, false},
{"missing pvc", true, func(b *BrokerConfig) { b.MetadataStorage.PvcSpec = nil }, false},
{"emptydir data", true, func(b *BrokerConfig) {
b.StorageConfigs = append(b.StorageConfigs, StorageConfig{MountPath: "/kafka-logs", EmptyDir: &corev1.EmptyDirVolumeSource{}})
}, false},
{"emptydir data without metadata storage", true, func(b *BrokerConfig) {
b.MetadataStorage = nil
b.StorageConfigs = []StorageConfig{{MountPath: "/kafka-logs", EmptyDir: &corev1.EmptyDirVolumeSource{}}}
}, true},
{"block pvc", true, func(b *BrokerConfig) {
mode := corev1.PersistentVolumeBlock
b.MetadataStorage.PvcSpec.VolumeMode = &mode
}, false},
{"reserved volume", true, func(b *BrokerConfig) { b.Volumes = []corev1.Volume{{Name: "kraft-metadata"}} }, false},
{"collision", true, func(b *BrokerConfig) { b.MetadataStorage.MountPath = "/csi-kafka-logs1" }, false},
{"nested collision", true, func(b *BrokerConfig) { b.MetadataStorage.MountPath = "/csi-kafka-logs1/metadata" }, false},
{"reserved mount", true, func(b *BrokerConfig) { b.MetadataStorage.MountPath = "/config/metadata" }, false},
{"custom mount", true, func(b *BrokerConfig) { b.VolumeMounts = []corev1.VolumeMount{{MountPath: "/csi-kafka-metadata/kafka"}} }, false},
{"relative path", true, func(b *BrokerConfig) { b.MetadataStorage.MountPath = "metadata" }, false},
{"unclean path", true, func(b *BrokerConfig) { b.MetadataStorage.MountPath = "/metadata/../logs" }, false},
} {
t.Run(tc.name, func(t *testing.T) {
b := base.DeepCopy()
tc.change(b)
if err := b.ValidateMetadataStorage(tc.kraft); (err == nil) != tc.valid {
t.Fatalf("valid=%v, got %v", tc.valid, err)
}
})
}
group := base.DeepCopy()
group.MetadataStorage.PvcSpec.StorageClassName = ptrString("group-class")
broker := Broker{BrokerConfigGroup: "brokers"}
spec := KafkaClusterSpec{BrokerConfigGroups: map[string]BrokerConfig{"brokers": *group}}
inherited, err := broker.GetBrokerConfig(spec)
if err != nil || inherited.MetadataStorage.MountPath != base.MetadataStorage.MountPath || len(inherited.StorageConfigs) != 1 {
t.Fatalf("inheritance failed: %v, %v", inherited, err)
}
broker.BrokerConfig = &BrokerConfig{MetadataStorage: &StorageConfig{MountPath: "/local-metadata", PvcSpec: &corev1.PersistentVolumeClaimSpec{}}}
local, err := broker.GetBrokerConfig(spec)
if err != nil || local.MetadataStorage.MountPath != "/local-metadata" || local.MetadataStorage.PvcSpec.StorageClassName != nil {
t.Fatalf("metadata override was merged rather than replaced: %v, %v", local, err)
}
local.MetadataStorage.MountPath = "/mutated"
if broker.BrokerConfig.MetadataStorage.MountPath != "/local-metadata" || group.MetadataStorage.MountPath != base.MetadataStorage.MountPath {
t.Fatal("merged config aliases its inputs")
}
a, b := base.MetadataStorage.DeepCopy(), base.MetadataStorage.DeepCopy()
b.PvcSpec.Resources.Requests = corev1.ResourceList{corev1.ResourceStorage: resource.MustParse("20Gi")}
if !MetadataStorageLocationEqual(a, b) {
t.Fatal("capacity change should not be a relocation")
}
b.MountPath = "/elsewhere"
if MetadataStorageLocationEqual(a, b) || MetadataStorageLocationEqual(a, nil) {
t.Fatal("removal/relocation allowed")
}
}

func ptrString(value string) *string { return &value }
5 changes: 5 additions & 0 deletions api/v1beta1/zz_generated.deepcopy.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading