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
80 changes: 80 additions & 0 deletions api/v1alpha1/hyperbytedbcluster_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,11 @@ type HyperbytedbClusterSpec struct {
// +optional
Cluster ClusterTuningSpec `json:"cluster,omitempty"`

// Experimental series sharding. When set, the operator writes a `[sharding]`
// block into config.toml. Enabling it requires replicas > 1 (cluster mode).
// +optional
Sharding *ShardingSpec `json:"sharding,omitempty"`

// +optional
Cardinality CardinalitySpec `json:"cardinality,omitempty"`

Expand Down Expand Up @@ -319,6 +324,81 @@ type ClusterTuningSpec struct {
DrainWaitSecs int32 `json:"drainWaitSecs,omitempty"`
}

// ShardingSpec maps to HyperbyteDB `[sharding]` (experimental series_id range
// sharding). Omitted fields are left out of config.toml so the server defaults apply.
// When enabled, config validation requires regionMergeSeries < regionSplitSeries
// < regionMaxSeries (same inequalities as hyperbytedb).
type ShardingSpec struct {
// Master switch. Written as sharding.enabled.
Enabled bool `json:"enabled"`

// Target replica count per shard region.
// +optional
// +kubebuilder:validation:Minimum=1
ReplicationFactor int32 `json:"replicationFactor,omitempty"`

// Target series per region before split.
// +optional
// +kubebuilder:validation:Minimum=1
RegionSplitSeries int64 `json:"regionSplitSeries,omitempty"`

// Hard split threshold (must be > regionSplitSeries when both are set).
// +optional
// +kubebuilder:validation:Minimum=1
RegionMaxSeries int64 `json:"regionMaxSeries,omitempty"`

// Merge when adjacent regions fall below this (must be < regionSplitSeries when both are set).
// +optional
// +kubebuilder:validation:Minimum=1
RegionMergeSeries int64 `json:"regionMergeSeries,omitempty"`

// Cooldown between split/merge on a region.
// +optional
// +kubebuilder:validation:Minimum=0
SplitMergeIntervalSecs int64 `json:"splitMergeIntervalSecs,omitempty"`

// Max concurrent split/move/merge operators.
// +optional
// +kubebuilder:validation:Minimum=1
ScheduleLimit int32 `json:"scheduleLimit,omitempty"`

// Region-stats report interval and shard scheduler tick (same duration).
// +optional
// +kubebuilder:validation:Minimum=1
HeartbeatIntervalSecs int64 `json:"heartbeatIntervalSecs,omitempty"`

// Sync bootstrap RPC timeout.
// +optional
// +kubebuilder:validation:Minimum=1
BootstrapTimeoutMs int64 `json:"bootstrapTimeoutMs,omitempty"`

// Seconds before the Raft leader proposes TransferPrimary for an unhealthy primary.
// +optional
// +kubebuilder:validation:Minimum=1
PrimaryFailoverAfterSecs int64 `json:"primaryFailoverAfterSecs,omitempty"`

// Per-peer HTTP timeout for sharded query/write/metadata scatter.
// +optional
// +kubebuilder:validation:Minimum=1
ScatterPeerTimeoutMs int64 `json:"scatterPeerTimeoutMs,omitempty"`

// Max Active peers tried per region per scatter request.
// +optional
// +kubebuilder:validation:Minimum=1
ScatterMaxPeerAttempts int32 `json:"scatterMaxPeerAttempts,omitempty"`

// Load-based split QPS threshold; 0 = disabled. Emitted only when set so the
// server default (0) applies when omitted.
// +optional
// +kubebuilder:validation:Minimum=0
LoadSplitQpsThreshold *int64 `json:"loadSplitQpsThreshold,omitempty"`

// Hard cap on regions per measurement.
// +optional
// +kubebuilder:validation:Minimum=1
MaxRegionsPerMeasurement int32 `json:"maxRegionsPerMeasurement,omitempty"`
}

// ReplicationSpec controls coordinator-side replication (how this node's
// accepted client writes are replicated to peers).
type ReplicationSpec struct {
Expand Down
23 changes: 23 additions & 0 deletions api/v1alpha1/hyperbytedbcluster_webhook.go
Original file line number Diff line number Diff line change
Expand Up @@ -150,5 +150,28 @@ func validateCluster(cluster *HyperbytedbCluster) (admission.Warnings, error) {
warnings = append(warnings, "2-node clusters cannot tolerate any node failure; consider 3+ replicas")
}

if err := validateSharding(cluster, replicas); err != nil {
return warnings, err
}

return warnings, nil
}

func validateSharding(cluster *HyperbytedbCluster, replicas int32) error {
s := cluster.Spec.Sharding
if s == nil {
return nil
}
if s.Enabled && replicas < 2 {
return fmt.Errorf("sharding.enabled requires replicas > 1 (cluster mode)")
}
if s.RegionMergeSeries > 0 && s.RegionSplitSeries > 0 && s.RegionMergeSeries >= s.RegionSplitSeries {
return fmt.Errorf("sharding.regionMergeSeries (%d) must be < regionSplitSeries (%d)",
s.RegionMergeSeries, s.RegionSplitSeries)
}
if s.RegionSplitSeries > 0 && s.RegionMaxSeries > 0 && s.RegionSplitSeries >= s.RegionMaxSeries {
return fmt.Errorf("sharding.regionSplitSeries (%d) must be < regionMaxSeries (%d)",
s.RegionSplitSeries, s.RegionMaxSeries)
}
return nil
}
89 changes: 89 additions & 0 deletions api/v1alpha1/hyperbytedbcluster_webhook_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
package v1alpha1

import (
"strings"
"testing"

"k8s.io/utils/ptr"
)

func TestValidateSharding(t *testing.T) {
t.Parallel()

tests := []struct {
name string
replicas int32
sharding *ShardingSpec
wantErr string
}{
{name: "nil sharding", replicas: 1},
{
name: "enabled requires cluster",
replicas: 1,
sharding: &ShardingSpec{Enabled: true},
wantErr: "replicas > 1",
},
{
name: "enabled on 3 replicas",
replicas: 3,
sharding: &ShardingSpec{Enabled: true, ReplicationFactor: 2},
},
{
name: "merge not less than split",
replicas: 3,
sharding: &ShardingSpec{
Enabled: true,
RegionMergeSeries: 10,
RegionSplitSeries: 5,
},
wantErr: "regionMergeSeries",
},
{
name: "split not less than max",
replicas: 3,
sharding: &ShardingSpec{
Enabled: true,
RegionSplitSeries: 10,
RegionMaxSeries: 10,
},
wantErr: "regionSplitSeries",
},
{
name: "kind split-test thresholds",
replicas: 6,
sharding: &ShardingSpec{
Enabled: true,
RegionSplitSeries: 5,
RegionMaxSeries: 10,
RegionMergeSeries: 2,
},
},
{
name: "disabled on single replica is ok",
replicas: 1,
sharding: &ShardingSpec{Enabled: false},
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
cluster := &HyperbytedbCluster{
Spec: HyperbytedbClusterSpec{
Replicas: ptr.To(tt.replicas),
Sharding: tt.sharding,
},
}
err := validateSharding(cluster, tt.replicas)
if tt.wantErr == "" {
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
return
}
if err == nil || !strings.Contains(err.Error(), tt.wantErr) {
t.Fatalf("want error containing %q, got %v", tt.wantErr, err)
}
})
}
}
25 changes: 25 additions & 0 deletions api/v1alpha1/zz_generated.deepcopy.go

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

Original file line number Diff line number Diff line change
Expand Up @@ -3589,6 +3589,89 @@ spec:
- enabled
type: object
type: object
sharding:
description: |-
Experimental series sharding. When set, the operator writes a `[sharding]`
block into config.toml. Enabling it requires replicas > 1 (cluster mode).
properties:
bootstrapTimeoutMs:
description: Sync bootstrap RPC timeout.
format: int64
minimum: 1
type: integer
enabled:
description: Master switch. Written as sharding.enabled.
type: boolean
heartbeatIntervalSecs:
description: Region-stats report interval and shard scheduler
tick (same duration).
format: int64
minimum: 1
type: integer
loadSplitQpsThreshold:
description: |-
Load-based split QPS threshold; 0 = disabled. Emitted only when set so the
server default (0) applies when omitted.
format: int64
minimum: 0
type: integer
maxRegionsPerMeasurement:
description: Hard cap on regions per measurement.
format: int32
minimum: 1
type: integer
primaryFailoverAfterSecs:
description: Seconds before the Raft leader proposes TransferPrimary
for an unhealthy primary.
format: int64
minimum: 1
type: integer
regionMaxSeries:
description: Hard split threshold (must be > regionSplitSeries
when both are set).
format: int64
minimum: 1
type: integer
regionMergeSeries:
description: Merge when adjacent regions fall below this (must
be < regionSplitSeries when both are set).
format: int64
minimum: 1
type: integer
regionSplitSeries:
description: Target series per region before split.
format: int64
minimum: 1
type: integer
replicationFactor:
description: Target replica count per shard region.
format: int32
minimum: 1
type: integer
scatterMaxPeerAttempts:
description: Max Active peers tried per region per scatter request.
format: int32
minimum: 1
type: integer
scatterPeerTimeoutMs:
description: Per-peer HTTP timeout for sharded query/write/metadata
scatter.
format: int64
minimum: 1
type: integer
scheduleLimit:
description: Max concurrent split/move/merge operators.
format: int32
minimum: 1
type: integer
splitMergeIntervalSecs:
description: Cooldown between split/merge on a region.
format: int64
minimum: 0
type: integer
required:
- enabled
type: object
statementSummary:
description: |-
StatementSummarySpec controls collection of per-statement execution stats
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,13 @@ spec:
replication:
mode: async
ackTimeoutMs: 5000
sharding:
enabled: true
replicationFactor: 3
regionSplitSeries: 100000
regionMaxSeries: 150000
regionMergeSeries: 20000
heartbeatIntervalSecs: 10
monitoring:
enabled: true
serviceMonitor: true
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,13 @@ spec:
ackTimeoutMs: 5000
syncQuorum:
minAcks: majority
sharding:
enabled: true
replicationFactor: 3
regionSplitSeries: 100000
regionMaxSeries: 150000
regionMergeSeries: 20000
heartbeatIntervalSecs: 10
monitoring:
enabled: true
serviceMonitor: true
Expand Down
Loading
Loading