diff --git a/apis/kubedb/v1/kafka_types.go b/apis/kubedb/v1/kafka_types.go index bf4f4f2cb2..43cd27f0f2 100644 --- a/apis/kubedb/v1/kafka_types.go +++ b/apis/kubedb/v1/kafka_types.go @@ -66,6 +66,17 @@ type KafkaSpec struct { // +optional Replicas *int32 `json:"replicas,omitempty"` + // Distributed if set true, manifestwork objects will be created instead of raw resources. + // A distributed Kafka is expanded by the operator into one self-contained Kafka cluster per + // Member data center for cross data center disaster recovery (DC-DR). + // +optional + Distributed bool `json:"distributed,omitempty"` + + // PodPlacementPolicy is the reference of the podPlacementPolicy that spreads the per data + // center Kafka clusters across data centers for DC-DR. + // +optional + PodPlacementPolicy *core.LocalObjectReference `json:"podPlacementPolicy,omitempty"` + // Kafka topology for node specification // +optional Topology *KafkaClusterTopology `json:"topology,omitempty"` @@ -184,6 +195,70 @@ type KafkaStatus struct { // Conditions applied to the database, such as approval or denial. // +optional Conditions []kmapi.Condition `json:"conditions,omitempty"` + // DisasterRecovery reports the cross data center (DC-DR) state for a distributed Kafka. + // +optional + DisasterRecovery *KafkaDisasterRecoveryStatus `json:"disasterRecovery,omitempty"` +} + +// KafkaDRPhase is the cross data center DR phase of a distributed Kafka. +type KafkaDRPhase string + +const ( + KafkaDRPhaseSteady KafkaDRPhase = "Steady" + KafkaDRPhaseFailingOver KafkaDRPhase = "FailingOver" + KafkaDRPhaseFailingBack KafkaDRPhase = "FailingBack" + KafkaDRPhaseDegraded KafkaDRPhase = "Degraded" +) + +// KafkaDisasterRecoveryStatus reports the per data center DC-DR view of a distributed Kafka. +// Kafka DC-DR is active/passive: exactly one data center is the write cluster, chosen by the +// dr-controlplane primary-DC Lease, and MirrorMaker 2 asynchronously mirrors it to the standby. +// This status reflects that decision on the single Kafka object. +type KafkaDisasterRecoveryStatus struct { + // ActiveDC is the data center that currently holds the primary DC Lease and takes producer writes. + // +optional + ActiveDC string `json:"activeDC,omitempty"` + + // Phase is the DC-DR phase. + // +optional + Phase KafkaDRPhase `json:"phase,omitempty"` + + // DataCenters is the per data center view, one entry per Member DC. + // +optional + DataCenters []KafkaDCStatus `json:"dataCenters,omitempty"` + + // LastTransitionTime is when ActiveDC last changed. + // +optional + LastTransitionTime *metav1.Time `json:"lastTransitionTime,omitempty"` +} + +// KafkaDCStatus is one data center's local view inside a distributed Kafka. +type KafkaDCStatus struct { + // ClusterName is the data center, named by its OCM managed cluster (the same + // clusterName used in the PlacementPolicy distributionRule). + ClusterName string `json:"clusterName"` + + // Role is Member or Arbiter. An Arbiter DC holds only the dr-controlplane etcd member and no Kafka. + // +optional + Role string `json:"role,omitempty"` + + // Writable is true when this DC is the active write cluster (its produce fence is open). + // +optional + Writable bool `json:"writable,omitempty"` + + // BrokersReady is the number of ready brokers in this DC's local Kafka cluster. + // +optional + BrokersReady *int32 `json:"brokersReady,omitempty"` + + // MirrorLagMillis is this DC's cross-DC MirrorMaker 2 replication lag behind the active DC, + // in milliseconds (the replication-latency-ms / record-age-ms metric, or the offset gap + // expressed as age). Nil when this DC is the active write cluster. + // +optional + MirrorLagMillis *int64 `json:"mirrorLagMillis,omitempty"` + + // Healthy reflects whether this DC's health Lease is fresh. + // +optional + Healthy bool `json:"healthy,omitempty"` } type KafkaTieredStorage struct { diff --git a/apis/kubedb/v1/zz_generated.deepcopy.go b/apis/kubedb/v1/zz_generated.deepcopy.go index 86e36b41fc..8604ab22c5 100644 --- a/apis/kubedb/v1/zz_generated.deepcopy.go +++ b/apis/kubedb/v1/zz_generated.deepcopy.go @@ -1006,6 +1006,11 @@ func (in *KafkaSpec) DeepCopyInto(out *KafkaSpec) { *out = new(int32) **out = **in } + if in.PodPlacementPolicy != nil { + in, out := &in.PodPlacementPolicy, &out.PodPlacementPolicy + *out = new(corev1.LocalObjectReference) + **out = **in + } if in.Topology != nil { in, out := &in.Topology, &out.Topology *out = new(KafkaClusterTopology) @@ -1093,6 +1098,11 @@ func (in *KafkaStatus) DeepCopyInto(out *KafkaStatus) { (*in)[i].DeepCopyInto(&(*out)[i]) } } + if in.DisasterRecovery != nil { + in, out := &in.DisasterRecovery, &out.DisasterRecovery + *out = new(KafkaDisasterRecoveryStatus) + (*in).DeepCopyInto(*out) + } return } @@ -1106,6 +1116,59 @@ func (in *KafkaStatus) DeepCopy() *KafkaStatus { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *KafkaDisasterRecoveryStatus) DeepCopyInto(out *KafkaDisasterRecoveryStatus) { + *out = *in + if in.DataCenters != nil { + in, out := &in.DataCenters, &out.DataCenters + *out = make([]KafkaDCStatus, len(*in)) + for i := range *in { + (*in)[i].DeepCopyInto(&(*out)[i]) + } + } + if in.LastTransitionTime != nil { + in, out := &in.LastTransitionTime, &out.LastTransitionTime + *out = (*in).DeepCopy() + } + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new KafkaDisasterRecoveryStatus. +func (in *KafkaDisasterRecoveryStatus) DeepCopy() *KafkaDisasterRecoveryStatus { + if in == nil { + return nil + } + out := new(KafkaDisasterRecoveryStatus) + in.DeepCopyInto(out) + return out +} + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *KafkaDCStatus) DeepCopyInto(out *KafkaDCStatus) { + *out = *in + if in.BrokersReady != nil { + in, out := &in.BrokersReady, &out.BrokersReady + *out = new(int32) + **out = **in + } + if in.MirrorLagMillis != nil { + in, out := &in.MirrorLagMillis, &out.MirrorLagMillis + *out = new(int64) + **out = **in + } + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new KafkaDCStatus. +func (in *KafkaDCStatus) DeepCopy() *KafkaDCStatus { + if in == nil { + return nil + } + out := new(KafkaDCStatus) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *KafkaTieredStorage) DeepCopyInto(out *KafkaTieredStorage) { *out = *in diff --git a/apis/ops/v1alpha1/kafka_ops_types.go b/apis/ops/v1alpha1/kafka_ops_types.go index bcbb5f3c8d..789fcec0e5 100644 --- a/apis/ops/v1alpha1/kafka_ops_types.go +++ b/apis/ops/v1alpha1/kafka_ops_types.go @@ -116,6 +116,23 @@ type KafkaHorizontalScalingSpec struct { Node *int32 `json:"node,omitempty"` // Node topology specification Topology *KafkaHorizontalScalingTopologySpec `json:"topology,omitempty"` + + // DataCenters scales individual data centers of a distributed DC-DR Kafka. + // Each entry sets that data center's local node count; data centers not listed + // are left unchanged. Use this instead of Node for a DC-DR cluster, where each + // data center runs its own self-contained Kafka cluster and is scaled independently. + // +optional + DataCenters []KafkaHorizontalScalingDC `json:"dataCenters,omitempty"` +} + +// KafkaHorizontalScalingDC is a per data center node-count target for scaling a +// distributed DC-DR Kafka. +type KafkaHorizontalScalingDC struct { + // ClusterName is the data center, named by its OCM managed cluster, matching a + // Member distributionRule in the Kafka PlacementPolicy. + ClusterName string `json:"clusterName"` + // Replicas is the desired local node count for this data center. + Replicas int32 `json:"replicas"` } // KafkaHorizontalScalingTopologySpec contains the horizontal scaling information in cluster topology mode diff --git a/apis/ops/v1alpha1/zz_generated.deepcopy.go b/apis/ops/v1alpha1/zz_generated.deepcopy.go index 361ac360da..d28c38858c 100644 --- a/apis/ops/v1alpha1/zz_generated.deepcopy.go +++ b/apis/ops/v1alpha1/zz_generated.deepcopy.go @@ -2867,6 +2867,11 @@ func (in *KafkaHorizontalScalingSpec) DeepCopyInto(out *KafkaHorizontalScalingSp *out = new(KafkaHorizontalScalingTopologySpec) (*in).DeepCopyInto(*out) } + if in.DataCenters != nil { + in, out := &in.DataCenters, &out.DataCenters + *out = make([]KafkaHorizontalScalingDC, len(*in)) + copy(*out, *in) + } return } @@ -2880,6 +2885,22 @@ func (in *KafkaHorizontalScalingSpec) DeepCopy() *KafkaHorizontalScalingSpec { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *KafkaHorizontalScalingDC) DeepCopyInto(out *KafkaHorizontalScalingDC) { + *out = *in + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new KafkaHorizontalScalingDC. +func (in *KafkaHorizontalScalingDC) DeepCopy() *KafkaHorizontalScalingDC { + if in == nil { + return nil + } + out := new(KafkaHorizontalScalingDC) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *KafkaHorizontalScalingTopologySpec) DeepCopyInto(out *KafkaHorizontalScalingTopologySpec) { *out = *in diff --git a/crds/kubedb.com_kafkas.yaml b/crds/kubedb.com_kafkas.yaml index f9195baa41..52980e4776 100644 --- a/crds/kubedb.com_kafkas.yaml +++ b/crds/kubedb.com_kafkas.yaml @@ -3367,6 +3367,8 @@ spec: type: string disableSecurity: type: boolean + distributed: + type: boolean enableSSL: type: boolean halted: @@ -3743,6 +3745,15 @@ spec: type: object type: object type: object + podPlacementPolicy: + default: + name: default + properties: + name: + default: "" + type: string + type: object + x-kubernetes-map-type: atomic podTemplate: properties: controller: diff --git a/crds/ops.kubedb.com_kafkaopsrequests.yaml b/crds/ops.kubedb.com_kafkaopsrequests.yaml index 7a78861532..d46c4bdd68 100644 --- a/crds/ops.kubedb.com_kafkaopsrequests.yaml +++ b/crds/ops.kubedb.com_kafkaopsrequests.yaml @@ -98,6 +98,19 @@ spec: x-kubernetes-map-type: atomic horizontalScaling: properties: + dataCenters: + items: + properties: + clusterName: + type: string + replicas: + format: int32 + type: integer + required: + - clusterName + - replicas + type: object + type: array node: format: int32 type: integer