【kubernetes v1.21】(kube-apiserver 1)kube-apiserver 核心架构与启动流程超深度分析
kube-apiserver 核心架构与启动流程超深度分析
基于 Kubernetes 源码(
cmd/kube-apiserver、pkg/controlplane、staging/src/k8s.io/apiserver)逐行拆解
一、模块定位
1.1 API Server 的业务职责
kube-apiserver 是 Kubernetes 控制平面的唯一入口,是整个集群的"大脑皮层"。它承担以下核心职责:
| 职责 | 说明 |
|---|---|
| REST API 网关 | 对外暴露 /api、/apis 等 HTTP 端点,是所有客户端(kubectl、controller-manager、scheduler、kubelet)的唯一交互入口 |
| 认证(Authentication) | 支持 X509客户端证书、Bearer Token、Bootstrap Token、OIDC、Webhook等多种认证方式,识别请求者身份 |
| 授权(Authorization) | 支持 RBAC、ABAC、Node、Webhook 等授权模式,判定请求者是否有权执行操作 |
| 准入控制(Admission) | 在对象持久化前执行深度检查与修改,包括 Mutating 和 Validating 两大类,支持内置插件与 Webhook |
| 数据持久化 | 将所有资源对象序列化后存入 etcd,是 etcd 的唯一直接消费者 |
| Watch 机制 | 基于 HTTP 长连接实现资源的变更通知,驱动控制器模式运行 |
| API 聚合(Aggregation) | 通过 kube-aggregator 将自定义 API Server(如 metrics-server)的 API 透明代理到统一入口 |
| CRD 管理 | 通过 apiextensions-apiserver 动态注册和管理 CustomResourceDefinition |
| 服务发现 | 提供 /apis 端点列举所有已注册的 API Group 和 Version |
| 安全策略 | 限流(MaxInFlight / PriorityAndFairness)、审计(Audit)、HSTS、CORS 等 |
1.2 在 Kubernetes 架构中的位置
┌─────────────────────────────────────────────────────┐
│ 外部客户端 │
│ (kubectl / SDK / curl / UI) │
└──────────────────────┬──────────────────────────────┘
│ HTTPS
▼
┌─────────────────────────────────────────────────────┐
│ ★ kube-apiserver ★ │
│ ┌──────────────────────────────────────────────┐ │
│ │ Aggregator Layer (kube-aggregator) │ │
│ │ ┌────────────────────────────────────────┐ │ │
│ │ │ APIExtensions Layer (CRD) │ │ │
│ │ │ ┌──────────────────────────────────┐ │ │ │
│ │ │ │ KubeAPIServer Core (内置资源) │ │ │ │
│ │ │ │ GenericAPIServer (通用框架) │ │ │ │
│ │ │ └──────────────────────────────────┘ │ │ │
│ │ └────────────────────────────────────────┘ │ │
│ └──────────────────────────────────────────────┘ │
└──────────────────────┬──────────────────────────────┘
│ gRPC / etcd v3 client
▼
┌─────────────────────────────────────────────────────┐
│ etcd 集群 │
└─────────────────────────────────────────────────────┘
kube-apiserver 处于所有组件的中心:
- 上行:所有客户端通过它读写集群状态
- 下行:它通过 etcd client 持久化数据
- 横向:controller-manager 和 scheduler 通过 Watch 机制获取变更通知
- 它是唯一直接与 etcd 通信的组件,保证了数据一致性
二、模块整体结构
2.1 Server Chain 三层委托结构
kube-apiserver 并非单一服务器,而是三层嵌套的委托链(Delegation Chain),每一层都是一个独立的 GenericAPIServer 实例:
Aggregator Server ──delegate──▶ KubeAPIServer ──delegate──▶ APIExtensions Server ──delegate──▶ EmptyDelegate
(最外层) (核心层) (CRD层) (链尾)
请求从最外层(Aggregator)进入,逐层向内委托处理:
- Aggregator:检查是否为已注册的 APIService,若是则代理到对应的后端 API Server
- APIExtensions:检查是否为 CRD 资源,若是则由 CRD handler 处理
- KubeAPIServer:处理所有 Kubernetes 内置资源(Pod、Service、Deployment 等)
- EmptyDelegate:返回 404,链尾哨兵
CreateServerChain 函数解析
func CreateServerChain(completedOptions completedServerRunOptions, stopCh <-chan struct{}) (*aggregatorapiserver.APIAggregator, error) {
// Step 1: 创建 NodeDialer(用于 SSH 隧道或代理连接 Kubelet)
nodeTunneler, proxyTransport, err := CreateNodeDialer(completedOptions)
// Step 2: 构建 KubeAPIServer 的 Config(最核心的配置)
kubeAPIServerConfig, serviceResolver, pluginInitializer, err := CreateKubeAPIServerConfig(...)
// Step 3: 构建 APIExtensions Server 的 Config
apiExtensionsConfig, err := createAPIExtensionsConfig(...)
// Step 4: 创建 APIExtensions Server(委托给 EmptyDelegate)
apiExtensionsServer, err := createAPIExtensionsServer(apiExtensionsConfig, genericapiserver.NewEmptyDelegate())
// Step 5: 创建 KubeAPIServer(委托给 APIExtensions Server)
kubeAPIServer, err := CreateKubeAPIServer(kubeAPIServerConfig, apiExtensionsServer.GenericAPIServer)
// Step 6: 构建 Aggregator Config
aggregatorConfig, err := createAggregatorConfig(...)
// Step 7: 创建 Aggregator Server(委托给 KubeAPIServer)
aggregatorServer, err := createAggregatorServer(aggregatorConfig, kubeAPIServer.GenericAPIServer, apiExtensionsServer.Informers)
return aggregatorServer, nil
}
关键设计原则:委托链从内到外构建,但请求从外到内流转。这种设计实现了关注点分离和可扩展性。
2.2 Config / CompletedConfig 全字段
genericapiserver.Config(通用服务器配置)
type Config struct {
// ===== 核心安全配置 =====
SecureServing *SecureServingInfo // TLS 服务配置
Authentication AuthenticationInfo // 认证配置
Authorization AuthorizationInfo // 授权配置
LoopbackClientConfig *restclient.Config // 回环客户端配置(用于 PostStartHook)
EgressSelector *egressselector.EgressSelector // 出站连接选择器
// ===== 访问控制 =====
RuleResolver authorizer.RuleResolver // 规则解析器
AdmissionControl admission.Interface // 准入控制器
FlowControl utilflowcontrol.Interface // 优先级与公平性限流
// ===== 服务能力 =====
CorsAllowedOriginList []string // CORS 允许的源
HSTSDirectives []string // HSTS 指令
EnableIndex bool // 启用 / 索引页
EnableProfiling bool // 启用 pprof
EnableDiscovery bool // 启用 API 发现
EnableContentionProfiling bool // 启用锁竞争分析
EnableMetrics bool // 启用 metrics 端点
// ===== 限流与超时 =====
MaxRequestsInFlight int // 最大非变更请求并发数(默认400)
MaxMutatingRequestsInFlight int // 最大变更请求并发数(默认200)
RequestTimeout time.Duration // 请求超时(默认60s)
MinRequestTimeout int // 最小请求超时(默认1800s,用于watch)
LivezGracePeriod time.Duration // 存活检查宽限期
ShutdownDelayDuration time.Duration // 关闭延迟时长
// ===== 请求限制 =====
JSONPatchMaxCopyBytes int64 // JSON Patch 最大复制字节数(默认3MB)
MaxRequestBodyBytes int64 // 最大请求体字节数(默认3MB)
GoawayChance float64 // HTTP/2 GOAWAY 概率(0~0.02)
LongRunningFunc apirequest.LongRunningRequestCheck // 长运行请求判断函数
// ===== Handler 构建 =====
BuildHandlerChainFunc func(apiHandler http.Handler, c *Config) http.Handler
HandlerChainWaitGroup *utilwaitgroup.SafeWaitGroup
// ===== API 资源 =====
Serializer runtime.NegotiatedSerializer // 序列化器
MergedResourceConfig *serverstore.ResourceConfig // 合并的资源启用配置
RESTOptionsGetter genericregistry.RESTOptionsGetter // REST 存储选项获取器
EquivalentResourceRegistry runtime.EquivalentResourceRegistry // 等价资源注册表
// ===== OpenAPI =====
OpenAPIConfig *openapicommon.Config // OpenAPI 规范配置
SkipOpenAPIInstallation bool // 跳过 OpenAPI 安装
// ===== 健康检查 =====
HealthzChecks []healthz.HealthChecker // healthz 检查
LivezChecks []healthz.HealthChecker // livez 检查
ReadyzChecks []healthz.HealthChecker // readyz 检查
// ===== Hook =====
PostStartHooks map[string]PostStartHookConfigEntry
DisabledPostStartHooks sets.String
// ===== 元数据 =====
Version *version.Info // 版本信息
ExternalAddress string // 外部访问地址
APIServerID string // API Server 唯一标识
StorageVersionManager storageversion.Manager // 存储版本管理器
DiscoveryAddresses discovery.Addresses // 发现地址
PublicAddress net.IP // 公网地址
LegacyAPIGroupPrefixes sets.String // 旧版 API 前缀
RequestInfoResolver apirequest.RequestInfoResolver // 请求信息解析器
// ===== 审计 =====
AuditBackend audit.Backend // 审计后端
AuditPolicyChecker auditpolicy.Checker // 审计策略检查器
}
controlplane.Config / ExtraConfig(KubeAPIServer 专属配置)
type Config struct {
GenericConfig *genericapiserver.Config // 通用配置
ExtraConfig ExtraConfig // KubeAPIServer 扩展配置
}
type ExtraConfig struct {
// ===== 集群认证 =====
ClusterAuthenticationInfo clusterauthenticationtrust.ClusterAuthenticationInfo
// ===== 存储与资源 =====
APIResourceConfigSource serverstorage.APIResourceConfigSource
StorageFactory serverserver.StorageFactory
EventTTL time.Duration
// ===== Kubelet 通信 =====
KubeletClientConfig kubeletclient.KubeletClientConfig
Tunneler tunneler.Tunneler
// ===== 服务网络 =====
ServiceIPRange net.IPNet // Service ClusterIP 范围
APIServerServiceIP net.IP // kubernetes Service 的 ClusterIP
SecondaryServiceIPRange net.IPNet // 双栈场景的二级 Service IP 范围
APIServerServicePort int // API Server Service 端口(默认443)
ServiceNodePortRange utilnet.PortRange // NodePort 范围
KubernetesServiceNodePort int // kubernetes Service 的 NodePort
// ===== Endpoint 协调 =====
MasterCount int // Master 数量
EndpointReconcilerType reconcilers.Type // Endpoint 协调器类型
EndpointReconcilerConfig EndpointReconcilerConfig
// ===== ServiceAccount =====
ServiceAccountIssuer serviceaccount.TokenGenerator
ServiceAccountMaxExpiration time.Duration
ExtendExpiration bool
ServiceAccountIssuerURL string
ServiceAccountJWKSURI string
ServiceAccountPublicKeys []interface{}
// ===== 其他 =====
EnableLogsSupport bool
ProxyTransport http.RoundTripper
VersionedInformers informers.SharedInformerFactory
IdentityLeaseDurationSeconds int
IdentityLeaseRenewIntervalSeconds int
}
completedConfig 模式
Kubernetes 使用 private struct + public wrapper 模式来强制"Complete→New"调用顺序:
// 私有结构体,外部无法直接构造
type completedConfig struct {
GenericConfig genericapiserver.CompletedConfig
ExtraConfig *ExtraConfig
}
// 公开包装,嵌入私有指针
type CompletedConfig struct {
*completedConfig // 外部只能通过 Config.Complete() 获得
}
// 同理,genericapiserver 也有相同模式:
type completedConfig struct { *Config; SharedInformerFactory informers.SharedInformerFactory }
type CompletedConfig struct { *completedConfig }
这种设计确保:你必须先调用 Complete() 填充默认值,然后才能用 New() 创建服务器实例——在编译期而非运行期强制约束。
2.3 ServerRunOptions 全字段
type ServerRunOptions struct {
// ===== 通用服务选项 =====
GenericServerRunOptions *genericoptions.ServerRunOptions
// 包含:AdvertiseAddress, CorsAllowedOriginList, HSTSDirectives, ExternalHost,
// MaxRequestsInFlight, MaxMutatingRequestsInFlight, RequestTimeout,
// GoawayChance, LivezGracePeriod, MinRequestTimeout, ShutdownDelayDuration,
// JSONPatchMaxCopyBytes, MaxRequestBodyBytes, EnablePriorityAndFairness
// ===== 存储选项 =====
Etcd *genericoptions.EtcdOptions
// 包含:StorageConfig(ServerList, Prefix, Codec等), EnableWatchCache, WatchCacheSizes等
// ===== 安全服务选项 =====
SecureServing *genericoptions.SecureServingOptionsWithLoopback
// 包含:BindAddress, BindPort, ServerCert, SNICertKeys等
// ===== 功能选项 =====
Features *genericoptions.FeatureOptions
// 包含:EnableWatchCache, APIEnablement等
// ===== 审计选项 =====
Audit *genericoptions.AuditOptions
// 包含:LogOptions, WebhookOptions, PolicyFile等
// ===== 准入控制选项 =====
Admission *kubeoptions.AdmissionOptions
// 包含:RecommendedAdmissionOptions(插件名列表、配置文件等)
// ===== 认证选项 =====
Authentication *kubeoptions.BuiltInAuthenticationOptions
// 包含:ClientCert, TokenFile, OIDC, ServiceAccounts, RequestHeader等
// ===== 授权选项 =====
Authorization *kubeoptions.BuiltInAuthorizationOptions
// 包含:Modes(RBAC,ABAC,Node,Webhook), PolicyFile, WebhookConfigFile等
// ===== 云提供商选项 =====
CloudProvider *kubeoptions.CloudProviderOptions
// ===== API 启用选项 =====
APIEnablement *genericoptions.APIEnablementOptions
// 包含:RuntimeConfig(map[string]bool)
// ===== 出站选择器 =====
EgressSelector *genericoptions.EgressSelectorOptions
// ===== 监控与日志 =====
Metrics *metrics.Options
Logs *logs.Options
// ===== KubeAPIServer 专属 =====
AllowPrivileged bool
EnableLogsHandler bool
EventTTL time.Duration // 默认 1h
KubeletConfig kubeletclient.KubeletClientConfig
KubernetesServiceNodePort int
MaxConnectionBytesPerSec int64
ServiceClusterIPRanges string // 用户输入
PrimaryServiceClusterIPRange net.IPNet // 解析后的主范围
SecondaryServiceClusterIPRange net.IPNet // 解析后的次范围
ServiceNodePortRange utilnet.PortRange // 默认 30000-32767
SSHKeyfile string // 已废弃
SSHUser string // 已废弃
ProxyClientCertFile string
ProxyClientKeyFile string
EnableAggregatorRouting bool
MasterCount int // 默认 1
EndpointReconcilerType string // 默认 "lease"
IdentityLeaseDurationSeconds int // 默认 3600
IdentityLeaseRenewIntervalSeconds int // 默认 10
ServiceAccountSigningKeyFile string
ServiceAccountIssuer serviceaccount.TokenGenerator
ServiceAccountTokenMaxExpiration time.Duration
ShowHiddenMetricsForVersion string
}
2.4 依赖注入关系
关键依赖注入点:
- StorageFactory →
EtcdOptions→Config.RESTOptionsGetter→ 各 RESTStorage - Authenticator →
AuthenticationOptions.ApplyTo()→Config.Authentication.Authenticator - Authorizer →
AuthorizationOptions.ToAuthorizationConfig().New()→Config.Authorization.Authorizer - AdmissionControl →
AdmissionOptions.ApplyTo()→Config.AdmissionControl - VersionedInformers →
LoopbackClientConfig→clientgoclientset→Config.ExtraConfig.VersionedInformers - FlowControl →
BuildPriorityAndFairness()→Config.FlowControl(需 APIPriorityAndFairness FeatureGate) - ServiceResolver →
buildServiceResolver()→AggregatorConfig.ExtraConfig.ServiceResolver
2.5 核心方法清单
| 方法 | 位置 | 作用 |
|---|---|---|
NewAPIServerCommand() | server.go | 创建 cobra.Command,定义 CLI 入口 |
Complete() | server.go | 填充默认值,返回 completedServerRunOptions |
Run() | server.go | 启动入口:CreateServerChain→PrepareRun→Run |
CreateServerChain() | server.go | 构建三层委托链 |
CreateNodeDialer() | server.go | 创建 SSH 隧道/代理传输 |
CreateKubeAPIServerConfig() | server.go | 构建 KubeAPIServer 配置 |
buildGenericConfig() | server.go | 构建通用配置(认证/授权/存储/准入等) |
BuildAuthorizer() | server.go | 构建授权器 |
BuildPriorityAndFairness() | server.go | 构建限流器 |
createAPIExtensionsConfig() | apiextensions.go | 构建 CRD Server 配置 |
createAPIExtensionsServer() | apiextensions.go | 创建 CRD Server |
createAggregatorConfig() | aggregator.go | 构建 Aggregator 配置 |
createAggregatorServer() | aggregator.go | 创建 Aggregator Server + 自动注册控制器 |
Config.Complete() | config.go | 完善通用配置 |
completedConfig.New() | config.go | 创建 GenericAPIServer 实例 |
GenericAPIServer.PrepareRun() | genericapiserver.go | 安装健康检查、OpenAPI |
preparedGenericAPIServer.Run() | genericapiserver.go | 启动 HTTP 服务 |
GenericAPIServer.InstallLegacyAPIGroup() | genericapiserver.go | 安装 /api/v1 资源 |
GenericAPIServer.InstallAPIGroups() | genericapiserver.go | 安装 /apis/* 资源 |
controlplane.Config.Complete() | instance.go | 完善 KubeAPIServer 配置 |
completedConfig.New() | instance.go | 创建 KubeAPIServer 实例+注册 PostStartHook |
2.6 数据流入流出
┌─────────────┐
│ 外部客户端 │──HTTPS──▶┌──────────────────────────────────────────────┐
│ (kubectl等) │ │ kube-apiserver │
└─────────────┘ │ │
│ 请求流向: │
│ Client ──▶ FullHandlerChain ──▶ Director │
│ │ │ │
│ │ Filter Chain: │ │
│ │ RequestReceivedTimestamp │ │
│ │ → PanicRecovery │ │
│ │ → RequestInfo │ │
│ │ → WaitGroup │ │
│ │ → RequestDeadline │ │
│ │ → TimeoutForNonLongRunning │ │
│ │ → CORS │ │
│ │ → Authentication ◄─── 认证信息 │ │
│ │ → Audit ◄───────── 审计策略 │ │
│ │ → Impersonation │ │
│ │ → MaxInFlight/PriorityAndFairness │ │
│ │ → Authorization ◄─── 授权信息 │ │
│ │ → StorageVersionPrecondition │ │
│ │ │ │
│ │ Director 路由: │ │
│ │ GoRestfulContainer ──▶ REST Storage │ │
│ │ NonGoRestfulMux ──▶ 其他 Handler │ │
│ │
│ 数据持久化: │
│ REST Storage ──▶ etcd3 ──▶ etcd 集群 │
│ │
│ Watch 通知: │
│ etcd Watch ──▶ Cacher ──▶ HTTP Chunked │
└──────────────────────────────────────────────┘
响应流向:
etcd ──▶ Decoder ──▶ Serializer ──▶ HTTP Response ──▶ Client
审计流向:
Request ──▶ Audit Policy Check ──▶ Audit Backend(Log/Webhook)
三、核心业务逻辑深度解析
3.1 main() → NewAPIServerCommand → Run → CreateServerChain 完整流程逐行解析
3.1.1 入口函数(隐式 main)
Kubernetes 的 cmd/kube-apiserver/apiserver.go 中:
func main() {
command := app.NewAPIServerCommand()
code := cli.Run(command)
os.Exit(code)
}
3.1.2 NewAPIServerCommand — 创建 Cobra 命令
func NewAPIServerCommand() *cobra.Command {
// 1. 创建默认 ServerRunOptions(所有配置的默认值在此确定)
s := options.NewServerRunOptions()
// 具体默认值:
// - MaxRequestsInFlight = 400
// - MaxMutatingRequestsInFlight = 200
// - RequestTimeout = 60s
// - MinRequestTimeout = 1800
// - EventTTL = 1h
// - MasterCount = 1
// - EndpointReconcilerType = "lease"
// - Etcd.DefaultStorageMediaType = "application/vnd.kubernetes.protobuf"
// - ServiceNodePortRange = 30000-32767
cmd := &cobra.Command{
Use: "kube-apiserver",
Long: `The Kubernetes API server validates and configures data...`,
SilenceUsage: true,
// PersistentPreRunE: 静默 client-go 警告(回环客户端不应产生自生警告)
PersistentPreRunE: func(*cobra.Command, []string) error {
rest.SetDefaultWarningHandler(rest.NoWarnings{})
return nil
},
// RunE: 核心执行逻辑
RunE: func(cmd *cobra.Command, args []string) error {
verflag.PrintAndExitIfRequested() // 处理 --version
cliflag.PrintFlags(fs) // 打印所有 flag 值(调试用)
checkNonZeroInsecurePort(fs) // 确保 insecure-port=0(已废弃)
// ★ Complete:填充默认值,转换选项
completedOptions, err := Complete(s)
// ★ Validate:校验所有选项
if errs := completedOptions.Validate(); len(errs) != 0 {
return utilerrors.NewAggregate(errs)
}
// ★ Run:启动服务器
return Run(completedOptions, genericapiserver.SetupSignalHandler())
// ↑ SetupSignalHandler() 返回 stopCh,
// 监听 SIGINT/SIGTERM 信号
},
}
// 注册所有 flag
fs := cmd.Flags()
namedFlagSets := s.Flags() // 分组注册:generic/etcd/secure serving/authentication/...
// 添加版本flag、全局flag、自定义全局flag
for _, f := range namedFlagSets.FlagSets {
fs.AddFlagSet(f)
}
return cmd
}
3.1.3 Complete — 选项补全
func Complete(s *options.ServerRunOptions) (completedServerRunOptions, error) {
// 1. 设置默认 AdvertiseAddress(若未指定,从 SecureServing 获取)
s.GenericServerRunOptions.DefaultAdvertiseAddress(s.SecureServing.SecureServingOptions)
// 2. 解析 ServiceClusterIPRanges → Primary + Secondary
// 支持双栈:如 "10.0.0.0/24,fd00::/108"
apiServerServiceIP, primaryRange, secondaryRange, err := getServiceIPAndRanges(s.ServiceClusterIPRanges)
s.PrimaryServiceClusterIPRange = primaryRange
s.SecondaryServiceClusterIPRange = secondaryRange
// 3. 若未提供 TLS 证书,自动生成自签名证书
s.SecureServing.MaybeDefaultWithSelfSignedCerts(
s.GenericServerRunOptions.AdvertiseAddress.String(),
[]string{"kubernetes.default.svc", "kubernetes.default", "kubernetes"},
[]net.IP{apiServerServiceIP},
)
// 4. 确定 ExternalHost
if len(s.GenericServerRunOptions.ExternalHost) == 0 {
if len(s.GenericServerRunOptions.AdvertiseAddress) > 0 {
s.GenericServerRunOptions.ExternalHost = s.GenericServerRunOptions.AdvertiseAddress.String()
} else {
hostname, _ := os.Hostname()
s.GenericServerRunOptions.ExternalHost = hostname
}
}
// 5. Authentication.ApplyAuthorization — 让认证选项知道授权模式
s.Authentication.ApplyAuthorization(s.Authorization)
// 6. ServiceAccount 签名密钥处理
// 若未设置 ServiceAccountSigningKeyFile,尝试使用 TLS 私钥
if s.ServiceAccountSigningKeyFile == "" {
if len(s.Authentication.ServiceAccounts.KeyFiles) == 0 && s.SecureServing.ServerCert.CertKey.KeyFile != "" {
if kubeauthenticator.IsValidServiceAccountKeyFile(s.SecureServing.ServerCert.CertKey.KeyFile) {
s.Authentication.ServiceAccounts.KeyFiles = []string{s.SecureServing.ServerCert.CertKey.KeyFile}
}
}
}
// 7. 构建 ServiceAccount Issuer(JWT Token Generator)
if s.ServiceAccountSigningKeyFile != "" && s.Authentication.ServiceAccounts.Issuer != "" {
sk, _ := keyutil.PrivateKeyFromFile(s.ServiceAccountSigningKeyFile)
// 校验 MaxExpiration 范围:[1h, 2^32s]
s.ServiceAccountIssuer, _ = serviceaccount.JWTTokenGenerator(s.Authentication.ServiceAccounts.Issuer, sk)
s.ServiceAccountTokenMaxExpiration = s.Authentication.ServiceAccounts.MaxExpiration
}
// 8. WatchCache 大小配置
if s.Etcd.EnableWatchCache {
sizes := kubeapiserver.DefaultWatchCacheSizes()
userSpecified, _ := serveroptions.ParseWatchCacheSizes(s.Etcd.WatchCacheSizes)
for resource, size := range userSpecified {
sizes[resource] = size // 用户覆盖默认值
}
s.Etcd.WatchCacheSizes, _ = serveroptions.WriteWatchCacheSizes(sizes)
}
// 9. 规范化 RuntimeConfig 中的 v1 前缀
// "v1" → "/v1", "api/v1" → "/v1", 删除 "api/legacy"
return completedServerRunOptions{ServerRunOptions: s}, nil
}
3.1.4 Run — 启动服务器
func Run(completeOptions completedServerRunOptions, stopCh <-chan struct{}) error {
klog.Infof("Version: %+v", version.Get()) // 打印版本信息
// ★ Step 1: 创建 Server Chain(三层委托链)
server, err := CreateServerChain(completeOptions, stopCh)
// ★ Step 2: 准备运行(安装健康检查、OpenAPI等)
prepared, err := server.PrepareRun()
// ★ Step 3: 运行服务器(阻塞直到 stopCh 关闭)
return prepared.Run(stopCh)
}
3.1.5 CreateServerChain — 构建三层委托链
这是最核心的函数,下面逐步拆解:
执行顺序:
1. CreateNodeDialer → 创建节点通信基础设施
2. CreateKubeAPIServerConfig → 构建核心配置
3. createAPIExtensionsConfig → 构建CRD配置
4. createAPIExtensionsServer → 创建CRD Server(delegate=EmptyDelegate)
5. CreateKubeAPIServer → 创建核心Server(delegate=CRD Server)
6. createAggregatorConfig → 构建聚合配置
7. createAggregatorServer → 创建聚合Server(delegate=KubeAPIServer)
Step 1: CreateNodeDialer — 创建节点拨号器
func CreateNodeDialer(s completedServerRunOptions) (tunneler.Tunneler, *http.Transport, error) {
var nodeTunneler tunneler.Tunneler
var proxyDialerFn utilnet.DialFunc
// 如果配置了 SSHUser,则建立 SSH 隧道
if len(s.SSHUser) > 0 {
// 初始化云提供商(获取 AddSSHKeyToAllInstances 方法)
cloud, _ := cloudprovider.InitCloudProvider(s.CloudProvider.CloudProvider, ...)
// 创建 Tunneler
nodeTunneler = tunneler.New(s.SSHUser, s.SSHKeyfile, healthCheckPath, installSSHKey)
proxyDialerFn = nodeTunneler.Dial
}
// 创建代理 Transport(InsecureSkipVerify=true,因为代理目标IP不可预知主机名)
proxyTransport := &http.Transport{
DialContext: proxyDialerFn,
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}
return nodeTunneler, proxyTransport, nil
}
Step 2: CreateKubeAPIServerConfig — 构建核心配置
func CreateKubeAPIServerConfig(s completedServerRunOptions, nodeTunneler tunneler.Tunneler, proxyTransport *http.Transport) (...) {
// ★ 核心调用:buildGenericConfig
genericConfig, versionedInformers, serviceResolver, pluginInitializers, admissionPostStartHook, storageFactory, err := buildGenericConfig(s.ServerRunOptions, proxyTransport)
// etcd 连接预检(重试60次,间隔1秒)
utilwait.PollImmediate(etcdRetryInterval, etcdRetryLimit*etcdRetryInterval,
preflight.EtcdConnection{ServerList: ...}.CheckEtcdServers)
// 初始化 capabilities(AllowPrivileged 等)
capabilities.Initialize(capabilities.Capabilities{...})
// 应用 metrics 和 logs 配置
s.Metrics.Apply()
s.Logs.Apply()
// 解析 Service IP 范围
serviceIPRange, apiServerServiceIP, _ := controlplane.ServiceIPRange(s.PrimaryServiceClusterIPRange)
// 构建 controlplane.Config
config := &controlplane.Config{
GenericConfig: genericConfig,
ExtraConfig: controlplane.ExtraConfig{
APIResourceConfigSource: storageFactory.APIResourceConfigSource,
StorageFactory: storageFactory,
EventTTL: s.EventTTL,
ServiceIPRange: serviceIPRange,
APIServerServiceIP: apiServerServiceIP,
// ... 其他字段
},
}
// 提取 ClientCA 和 RequestHeader 配置到 ClusterAuthenticationInfo
clientCAProvider, _ := s.Authentication.ClientCert.GetClientCAContentProvider()
config.ExtraConfig.ClusterAuthenticationInfo.ClientCA = clientCAProvider
requestHeaderConfig, _ := s.Authentication.RequestHeader.ToAuthenticationRequestHeaderConfig()
config.ExtraConfig.ClusterAuthenticationInfo.RequestHeaderCA = requestHeaderConfig.CAContentProvider
// ... 其他 RequestHeader 字段
// 添加准入初始化 PostStartHook
config.GenericConfig.AddPostStartHook("start-kube-apiserver-admission-initializer", admissionPostStartHook)
// 处理 EgressSelector(出站流量控制)
if config.GenericConfig.EgressSelector != nil {
config.ExtraConfig.KubeletClientConfig.Lookup = config.GenericConfig.EgressSelector.Lookup
// 修改 ProxyTransport 的 DialContext
}
// 加载 ServiceAccount 公钥
var pubKeys []interface{}
for _, f := range s.Authentication.ServiceAccounts.KeyFiles {
keys, _ := keyutil.PublicKeysFromFile(f)
pubKeys = append(pubKeys, keys...)
}
config.ExtraConfig.ServiceAccountPublicKeys = pubKeys
return config, serviceResolver, pluginInitializers, nil
}
buildGenericConfig — 最核心的配置构建函数:
func buildGenericConfig(s *options.ServerRunOptions, proxyTransport *http.Transport) (...) {
// 1. 创建通用 Config
genericConfig = genericapiserver.NewConfig(legacyscheme.Codecs)
genericConfig.MergedResourceConfig = controlplane.DefaultAPIResourceConfigSource()
// 2. 依次 Apply 各 Options
s.GenericServerRunOptions.ApplyTo(genericConfig) // 地址、限流、超时
s.SecureServing.ApplyTo(&genericConfig.SecureServing, &genericConfig.LoopbackClientConfig)
s.Features.ApplyTo(genericConfig) // WatchCache
s.APIEnablement.ApplyTo(genericConfig, ...) // 资源启用
s.EgressSelector.ApplyTo(genericConfig) // 出站选择
// 3. 设置 OpenAPI 配置
genericConfig.OpenAPIConfig = genericapiserver.DefaultOpenAPIConfig(
generatedopenapi.GetOpenAPIDefinitions,
openapinamer.NewDefinitionNamer(legacyscheme.Scheme, extensionsapiserver.Scheme, aggregatorscheme.Scheme),
)
genericConfig.OpenAPIConfig.Info.Title = "Kubernetes"
// 4. 设置长运行请求判断函数
genericConfig.LongRunningFunc = filters.BasicLongRunningRequestCheck(
sets.NewString("watch", "proxy"), // 长运行 verbs
sets.NewString("attach", "exec", "proxy", "log", "portforward"), // 长运行 subresources
)
// 5. 构建 StorageFactory
storageFactoryConfig := kubeapiserver.NewStorageFactoryConfig()
completedStorageFactoryConfig, _ := storageFactoryConfig.Complete(s.Etcd)
storageFactory, _ = completedStorageFactoryConfig.New()
s.Etcd.ApplyWithStorageFactoryTo(storageFactory, genericConfig)
// 6. 配置回环客户端(用 protobuf 自通信,禁用压缩)
genericConfig.LoopbackClientConfig.ContentConfig.ContentType = "application/vnd.kubernetes.protobuf"
genericConfig.LoopbackClientConfig.DisableCompression = true
// 7. 创建外部 clientset 和 SharedInformerFactory
clientgoExternalClient, _ := clientgoclientset.NewForConfig(genericConfig.LoopbackClientConfig)
versionedInformers = clientgoinformers.NewSharedInformerFactory(clientgoExternalClient, 10*time.Minute)
// 8. 认证配置
s.Authentication.ApplyTo(&genericConfig.Authentication, genericConfig.SecureServing,
genericConfig.EgressSelector, genericConfig.OpenAPIConfig, clientgoExternalClient, versionedInformers)
// 9. 授权配置
genericConfig.Authorization.Authorizer, genericConfig.RuleResolver, _ = BuildAuthorizer(s, genericConfig.EgressSelector, versionedInformers)
// 10. 若未启用 RBAC,禁用 RBAC PostStartHook
if !sets.NewString(s.Authorization.Modes...).Has(modes.ModeRBAC) {
genericConfig.DisabledPostStartHooks.Insert(rbacrest.PostStartHookName)
}
// 11. 审计配置
s.Audit.ApplyTo(genericConfig)
// 12. 准入控制配置
admissionConfig := &kubeapiserveradmission.Config{
ExternalInformers: versionedInformers,
LoopbackClientConfig: genericConfig.LoopbackClientConfig,
CloudConfigFile: s.CloudProvider.CloudConfigFile,
}
serviceResolver = buildServiceResolver(s.EnableAggregatorRouting, ...)
pluginInitializers, admissionPostStartHook, _ = admissionConfig.New(proxyTransport, genericConfig.EgressSelector, serviceResolver)
s.Admission.ApplyTo(genericConfig, versionedInformers, kubeClientConfig, feature.DefaultFeatureGate, pluginInitializers...)
// 13. 优先级与公平性限流
if utilfeature.DefaultFeatureGate.Enabled(genericfeatures.APIPriorityAndFairness) && s.GenericServerRunOptions.EnablePriorityAndFairness {
genericConfig.FlowControl = BuildPriorityAndFairness(s, clientgoExternalClient, versionedInformers)
}
return
}
3.2 CreateServerChain 三层委托链深度解析
3.2.1 APIExtensions Server 构建
func createAPIExtensionsConfig(kubeAPIServerConfig genericapiserver.Config, ...) (*apiextensionsapiserver.Config, error) {
// ★ 关键:浅拷贝 genericConfig
genericConfig := kubeAPIServerConfig
// 清空 PostStartHooks(避免重复注册)
genericConfig.PostStartHooks = map[string]genericapiserver.PostStartHookConfigEntry{}
// 清空 RESTOptionsGetter(后面用自己的)
genericConfig.RESTOptionsGetter = nil
// 重新应用 Admission(使用 apiextensions 自己的 scheme)
commandOptions.Admission.ApplyTo(&genericConfig, ...)
// 覆盖 Etcd 配置:使用 apiextensions 的 Codec
etcdOptions := *commandOptions.Etcd
etcdOptions.StorageConfig.Codec = apiextensionsapiserver.Codecs.LegacyCodec(v1beta1, v1)
etcdOptions.StorageConfig.EncodeVersioner = runtime.NewMultiGroupVersioner(v1beta1.SchemeGroupVersion, ...)
genericConfig.RESTOptionsGetter = &genericoptions.SimpleRestOptionsFactory{Options: etcdOptions}
// 覆盖 MergedResourceConfig
commandOptions.APIEnablement.ApplyTo(&genericConfig,
apiextensionsapiserver.DefaultAPIResourceConfigSource(),
apiextensionsapiserver.Scheme)
return &apiextensionsapiserver.Config{
GenericConfig: &genericapiserver.RecommendedConfig{
Config: genericConfig,
SharedInformerFactory: externalInformers,
},
ExtraConfig: apiextensionsapiserver.ExtraConfig{
CRDRESTOptionsGetter: ...,
MasterCount: masterCount,
AuthResolverWrapper: authResolverWrapper,
ServiceResolver: serviceResolver,
},
}, nil
}
3.2.2 KubeAPIServer 构建
func CreateKubeAPIServer(kubeAPIServerConfig *controlplane.Config, delegateAPIServer genericapiserver.DelegationTarget) (*controlplane.Instance, error) {
// Config.Complete() → CompletedConfig
// CompletedConfig.New(delegate) → Instance
return kubeAPIServerConfig.Complete().New(delegateAPIServer)
}
completedConfig.New() 做了以下事情:
- 调用
GenericConfig.New("kube-apiserver", delegate)创建GenericAPIServer - 安装 Logs 路由(如启用)
- 安装 OpenID 元数据端点
- 安装 Legacy API(/api/v1):Pod、Service、Node、ConfigMap 等
- 安装 API Groups(/apis/*):apps、batch、rbac 等 18 个组
- 安装 Tunneler(如有)
- 注册 PostStartHook:
start-cluster-authentication-info-controller - 注册 PostStartHook:
start-kube-apiserver-identity-lease-controller(如 APIServerIdentity FeatureGate 启用) - 注册 PostStartHook:
start-kube-apiserver-identity-lease-garbage-collector
3.2.3 Aggregator Server 构建
func createAggregatorServer(aggregatorConfig *aggregatorapiserver.Config, delegateAPIServer genericapiserver.DelegationTarget, ...) (*aggregatorapiserver.APIAggregator, error) {
// 创建 Aggregator Server
aggregatorServer, _ := aggregatorConfig.Complete().NewWithDelegate(delegateAPIServer)
// 创建自动注册控制器
autoRegistrationController := autoregister.NewAutoRegisterController(...)
apiServices := apiServicesToRegister(delegateAPIServer, autoRegistrationController)
// 创建 CRD 注册控制器
crdRegistrationController := crdregistration.NewCRDRegistrationController(...)
// 注册 PostStartHook:自动注册 APIService
aggregatorServer.GenericAPIServer.AddPostStartHook("kube-apiserver-autoregistration", func(context) error {
go crdRegistrationController.Run(5, context.StopCh)
go func() {
// 等待 CRD 初始同步完成后再启动自动注册
if aggregatorConfig.GenericConfig.MergedResourceConfig.AnyVersionForGroupEnabled("apiextensions.k8s.io") {
crdRegistrationController.WaitForInitialSync()
}
autoRegistrationController.Run(5, context.StopCh)
}()
return nil
})
// 添加 APIService 可用性健康检查
aggregatorServer.GenericAPIServer.AddBootSequenceHealthChecks(
makeAPIServiceAvailableHealthCheck("autoregister-completion", apiServices, ...))
return aggregatorServer, nil
}
3.3 GenericServer 构建过程
func (c completedConfig) New(name string, delegationTarget DelegationTarget) (*GenericAPIServer, error) {
// 前置检查
if c.Serializer == nil { return nil, fmt.Errorf("Serializer == nil") }
if c.LoopbackClientConfig == nil { return nil, fmt.Errorf("LoopbackClientConfig == nil") }
if c.EquivalentResourceRegistry == nil { return nil, fmt.Errorf("EquivalentResourceRegistry == nil") }
// ★ 构建 Handler Chain Builder
handlerChainBuilder := func(handler http.Handler) http.Handler {
return c.BuildHandlerChainFunc(handler, c.Config)
}
// ★ 创建 APIServerHandler(核心 HTTP 路由器)
apiServerHandler := NewAPIServerHandler(name, c.Serializer, handlerChainBuilder, delegationTarget.UnprotectedHandler())
// APIServerHandler 包含:
// - FullHandlerChain: handlerChainBuilder(director) = 完整过滤链
// - GoRestfulContainer: go-restful 容器(API 路由)
// - NonGoRestfulMux: PathRecorderMux(非 API 路由)
// - Director: 路由分发器
// ★ 创建 GenericAPIServer
s := &GenericAPIServer{
discoveryAddresses: c.DiscoveryAddresses,
LoopbackClientConfig: c.LoopbackClientConfig,
legacyAPIGroupPrefixes: c.LegacyAPIGroupPrefixes,
admissionControl: c.AdmissionControl,
Serializer: c.Serializer,
AuditBackend: c.AuditBackend,
Authorizer: c.Authorization.Authorizer,
delegationTarget: delegationTarget,
Handler: apiServerHandler,
// Hook 管理
postStartHooks: map[string]postStartHookEntry{},
preShutdownHooks: map[string]preShutdownHookEntry{},
disabledPostStartHooks: c.DisabledPostStartHooks,
// 健康检查
healthzChecks: c.HealthzChecks,
livezChecks: c.LivezChecks,
readyzChecks: c.ReadyzChecks,
readinessStopCh: make(chan struct{}),
// 其他
DiscoveryGroupManager: discovery.NewRootAPIsHandler(c.DiscoveryAddresses, c.Serializer),
minRequestTimeout: time.Duration(c.MinRequestTimeout) * time.Second,
ShutdownTimeout: c.RequestTimeout,
ShutdownDelayDuration: c.ShutdownDelayDuration,
SecureServingInfo: c.SecureServing,
ExternalAddress: c.ExternalAddress,
APIServerID: c.APIServerID,
StorageVersionManager: c.StorageVersionManager,
}
// 设置 JSON Patch 大小限制
atomic.CompareAndSwapInt64(&jsonpatch.AccumulatedCopySizeLimit, existing, c.JSONPatchMaxCopyBytes)
// ★ 从委托目标继承 PostStartHooks
for k, v := range delegationTarget.PostStartHooks() {
s.postStartHooks[k] = v
}
for k, v := range delegationTarget.PreShutdownHooks() {
s.preShutdownHooks[k] = v
}
// ★ 注册预配置的 PostStartHooks
for name, preconfiguredPostStartHook := range c.PostStartHooks {
s.AddPostStartHook(name, preconfiguredPostStartHook.hook)
}
// ★ 注册内置 PostStartHooks
// 1. generic-apiserver-start-informers: 启动 SharedInformerFactory
s.AddPostStartHook("generic-apiserver-start-informers", func(context) error {
c.SharedInformerFactory.Start(context.StopCh)
return nil
})
// 2. priority-and-fairness-config-consumer: 启动限流控制器
if c.FlowControl != nil {
s.AddPostStartHook("priority-and-fairness-config-consumer", func(context) error {
go c.FlowControl.MaintainObservations(context.StopCh)
go c.FlowControl.Run(context.StopCh)
return nil
})
// 3. priority-and-fairness-filter: 启动水位线维护
s.AddPostStartHook("priority-and-fairness-filter", func(context) error {
genericfilters.StartPriorityAndFairnessWatermarkMaintenance(context.StopCh)
return nil
})
} else {
// 3b. max-in-flight-filter: 启动最大并发水位线维护
s.AddPostStartHook("max-in-flight-filter", ...)
}
// ★ 继承委托目标的健康检查
for _, delegateCheck := range delegationTarget.HealthzChecks() {
s.AddHealthChecks(delegateCheck)
}
// ★ 合并 listedPathProvider
s.listedPathProvider = routes.ListedPathProviders{s.listedPathProvider, delegationTarget}
// ★ 安装基础 API 路由
installAPI(s, c.Config)
// 包括:Index、Profiling、Metrics、Version、Discovery、FlowControl
return s, nil
}
3.4 PostStartHook 机制
PostStartHook 是 kube-apiserver 的延迟初始化机制:在服务器开始监听端口后,以独立 goroutine 并发执行注册的 Hook 函数。
3.4.1 Hook 注册
type PostStartHookFunc func(context PostStartHookContext) error
type PostStartHookContext struct {
LoopbackClientConfig *restclient.Config // 回环客户端(特权访问)
StopCh <-chan struct{} // 停止信号
}
type postStartHookEntry struct {
hook PostStartHookFunc
originatingStack string // 调试用:注册时的调用栈
done chan struct{} // 完成信号(用于健康检查)
}
注册流程:
- 检查名称非空、hook 非 nil
- 检查是否被 DisabledPostStartHooks 禁用
- 检查是否已被调用过(
postStartHooksCalled标志) - 检查是否重名
- 创建
donechannel - 注册
postStartHookHealthz(检查 hook 是否完成的健康检查) - 存入
postStartHooksmap
3.4.2 Hook 执行
func (s *GenericAPIServer) RunPostStartHooks(stopCh <-chan struct{}) {
s.postStartHookLock.Lock()
defer s.postStartHookLock.Unlock()
s.postStartHooksCalled = true // ★ 标记已调用,后续注册将被拒绝
context := PostStartHookContext{
LoopbackClientConfig: s.LoopbackClientConfig,
StopCh: stopCh,
}
// ★ 每个 Hook 独立 goroutine,并发执行,无顺序保证
for hookName, hookEntry := range s.postStartHooks {
go runPostStartHook(hookName, hookEntry, context)
}
}
func runPostStartHook(name string, entry postStartHookEntry, context PostStartHookContext) {
defer utilruntime.HandleCrash() // ★ 防止意外 panic 杀死服务器
err := entry.hook(context)
if err != nil {
klog.Fatalf("PostStartHook %q failed: %v", name, err) // ★ 故意杀死:Hook 失败=服务器不可用
}
close(entry.done) // ★ 通知健康检查:Hook 已完成
}
3.4.3 KubeAPIServer 中注册的所有 PostStartHooks
| Hook 名称 | 作用 |
|---|---|
start-kube-apiserver-admission-initializer | 初始化准入插件 |
generic-apiserver-start-informers | 启动 SharedInformerFactory |
priority-and-fairness-config-consumer | 启动 APF 限流控制器 |
priority-and-fairness-filter | 启动 APF 水位线维护 |
max-in-flight-filter | 启动最大并发水位线维护(无APF时) |
bootstrap-controller | 启动 Bootstrap 控制器(kubernetes Service/Endpoint) |
start-cluster-authentication-info-controller | 同步 ClusterAuthenticationInfo 到 ConfigMap |
start-kube-apiserver-identity-lease-controller | 维护 API Server 身份 Lease |
start-kube-apiserver-identity-lease-garbage-collector | GC 过期的身份 Lease |
kube-apiserver-autoregistration | 自动注册 APIService + CRD 注册控制器 |
start-kube-apiserver-informers | 各组 Informer 的启动 |
| RBAC related hooks | RBAC 数据引导等 |
3.5 优雅关闭流程
关键时序:
- 收到信号 →
stopCh关闭 - 立即关闭
readinessStopCh→/readyz返回失败,负载均衡器开始摘除 - 等待
ShutdownDelayDuration→ 给 LB 时间完成摘除(不再有新流量) - 关闭
delayedStopCh→ 触发 HTTP Server 优雅关闭 - 执行 PreShutdownHooks → 清理操作(如从 kubernetes Endpoint 中移除自己)
- 等待 HTTP Server 关闭完成 →
server.Shutdown()等待所有活跃请求完成 - 等待 HandlerChainWaitGroup → 确保所有过滤链中的请求完成
- 完成
3.6 SecureServing / Authentication / Authorization 配置构建
3.6.1 SecureServing
type SecureServingInfo struct {
Listener net.Listener // 网络监听器
Cert dynamiccertificates.CertKeyContentProvider // 主证书
SNICerts []dynamiccertificates.SNICertKeyContentProvider // SNI 证书
ClientCA dynamiccertificates.CAContentProvider // 客户端 CA
MinTLSVersion uint16 // 最低 TLS 版本
CipherSuites []uint16 // 允许的密码套件
HTTP2MaxStreamsPerConnection int // HTTP/2 最大并发流
DisableHTTP2 bool // 禁用 HTTP/2
}
SecureServingInfo.Serve() 做了以下事情:
- 构建
tls.Config(最低 TLS 1.2,启用 HTTP/2) - 若有 ClientCA:设置
tls.RequestClientCert(请求但不强制客户端证书) - 创建
DynamicServingCertificateController(动态证书热更新) - 注册证书变更监听器
- 启动证书控制器 goroutine
- 配置 HTTP/2 参数(MaxUploadBufferPerStream=256KB, MaxConcurrentStreams=250)
- 创建
http.Server并启动RunServer()
3.6.2 Authentication
认证通过 BuiltInAuthenticationOptions.ApplyTo() 构建:
// 支持的认证方式(按优先级合并为 union authenticator):
// 1. Loopback Token(回环客户端 Bearer Token,最高优先级)
// 2. X509 Client Certificate
// 3. Bearer Token(静态 TokenFile)
// 4. Bootstrap Token
// 5. OIDC
// 6. ServiceAccount Token
// 7. RequestHeader(代理认证)
// 8. Webhook Token Authenticator
认证流程:
- 每种认证方式创建对应的
Authenticator实现 - 通过
authenticatorunion.New()合并为 Union Authenticator - 设置到
Config.Authentication.Authenticator - 同时将 Loopback Token 注册为特权认证和授权(
AuthorizeClientBearerToken)
3.6.3 Authorization
func BuildAuthorizer(s *options.ServerRunOptions, EgressSelector, versionedInformers) (authorizer.Authorizer, authorizer.RuleResolver, error) {
authorizationConfig := s.Authorization.ToAuthorizationConfig(versionedInformers)
// 若有 EgressSelector,设置自定义拨号器
if EgressSelector != nil {
egressDialer, _ := EgressSelector.Lookup(egressselector.ControlPlane.AsNetworkContext())
authorizationConfig.CustomDial = egressDialer
}
return authorizationConfig.New()
// 创建 Union Authorizer,包含:
// - RBAC Authorizer(基于 Role/ClusterRole/RoleBinding/ClusterRoleBinding)
// - Node Authorizer(基于 Node 的权限限制)
// - ABAC Authorizer(基于属性的访问控制)
// - Webhook Authorizer(外部授权服务)
}
四、Mermaid 图集
图1:整体架构图
图2:启动完整流程图
图3:Server 委托链图
图4:核心类图
图5:依赖注入图
图6:Config 构建图
图7:PostStartHook 机制图
图8:优雅关闭图
图9:数据流入流出图
图10:接口层次图
图11:Options 层次图
图12:ServerChain 构建时序图
图13:健康检查图
图14:Delegate Chain 详细路由图
图15:路由注册图
五、总结与关键洞察
5.1 设计模式总结
| 模式 | 应用 |
|---|---|
| Delegation Chain | 三层服务器通过委托链组合,实现关注点分离与可扩展性 |
| CompletedConfig Pattern | 私有结构体+公开包装,强制 Complete→New 调用顺序 |
| Builder Pattern | Options→ApplyTo→Config→New,逐步构建复杂配置 |
| Hook Pattern | PostStartHook/PreShutdownHook 实现延迟初始化与优雅清理 |
| Strategy Pattern | Authenticator/Authorizer/Admission 通过接口替换具体实现 |
| Observer Pattern | dynamiccertificates.Notifier/Listener 实现证书热更新 |
| Factory Pattern | StorageFactory、RESTStorageProvider 封装资源创建逻辑 |
5.2 关键架构决策
-
Shallow Copy 配置传递:APIExtensions 和 Aggregator 对 KubeAPIServer 的 Config 进行浅拷贝,共享大部分配置但覆盖关键差异(Codec、RESTOptionsGetter、Scheme)
-
Protobuf 自通信:LoopbackClientConfig 使用 protobuf 编码(最高效),禁用压缩(本地通信无需压缩),这是 apiserver 内部通信的性能优化
-
PostStartHook 的容错设计:每个 Hook 有独立的
donechannel 和健康检查,Hook 意外 panic 被HandleCrash()捕获,但故意失败(klog.Fatalf)会杀死进程——因为 Hook 失败意味着服务器不可用 -
优雅关闭的分层设计:先标记不就绪(/readyz 失败)→ 等待 LB 摘除(ShutdownDelayDuration)→ 执行 PreShutdownHook → 关闭 HTTP Server → 等待请求完成
-
Handler Chain 的洋葱模型:过滤链从外到内依次包裹,每个 Filter 负责一个关注点,最内层是 Director 路由分发器
5.3 启动时序关键点
- Options → Config 是单向的、不可逆的转化
- Server Chain 从内到外构建:EmptyDelegate → APIExtensions → KubeAPIServer → Aggregator
- 请求从外到内流转:Aggregator → KubeAPIServer → APIExtensions → EmptyDelegate
- PostStartHook 在 HTTP 服务启动后执行,否则无法使用 LoopbackClientConfig
- systemd 通知在所有 Hook 启动后发送,表示服务就绪
本文档基于 Kubernetes 源码
cmd/kube-apiserver/app/server.go、apiextensions.go、aggregator.go、pkg/controlplane/instance.go、staging/src/k8s.io/apiserver/pkg/server/config.go、genericapiserver.go、handler.go、hooks.go、secure_serving.go等核心文件深度分析编写。
更多推荐


所有评论(0)