【kubernetes v1.21】(二)kube-apiserver 超深度架构分析
kube-apiserver 超深度架构分析
基于 Kubernetes 源码逐行分析,目标模块:kube-apiserver 及其关联组件
一、模块定位
1.1 业务职责
kube-apiserver 是 Kubernetes 控制面的核心组件,是整个集群的唯一入口(API Gateway)。其职责涵盖:
- REST API 服务:暴露集群管理 API,所有组件(kubelet、kube-proxy、controller-manager、scheduler)及外部客户端均通过它与集群交互
- 认证(Authentication):验证请求者身份,支持 X509证书、Bearer Token、OIDC、Webhook、Bootstrap Token 等多种机制
- 授权(Authorization):基于 RBAC/ABAC/Webhook/Node 等模式判定请求是否被允许
- 准入控制(Admission Control):在对象持久化前进行变更(Mutating)和验证(Validating),包含内置插件和 Webhook 扩展
- 数据持久化:将 API 对象通过 etcd v3 存储后端持久化,支持 Watch 机制
- API 聚合(Aggregation):通过 kube-aggregator 将扩展 API Server 注册到统一入口
- CRD 支持:通过 apiextensions-apiserver 动态注册自定义资源
- 审计(Audit):记录请求级别的审计日志
- API 优先级与公平性(API Priority and Fairness):基于 FlowControl 对请求进行限流和调度
1.2 在系统中的位置
kube-apiserver 位于 Kubernetes 架构的中心枢纽位置:
Client (kubectl/SDK) ──HTTPS──▶ kube-apiserver ──gRPC──▶ etcd
│
┌────────────────────┼────────────────────┐
│ │ │
kubelet controller-manager scheduler
kube-apiserver 本身是一个三层委托链(Delegation Chain)的聚合体:
请求 ──▶ Aggregator (kube-aggregator)
│ (未命中则委托)
▼
KubeAPIServer (核心内置 API)
│ (未命中则委托)
▼
APIExtensions (CRD API)
│ (未命中则委托)
▼
EmptyDelegate (404)
二、模块整体结构
2.1 代码目录与文件布局
| 目录 | 职责 |
|---|---|
cmd/kube-apiserver/apiserver.go | 进程入口 main() |
cmd/kube-apiserver/app/server.go | 命令构造、启动流程、服务器链构建 |
cmd/kube-apiserver/app/aggregator.go | Aggregator 配置与创建 |
cmd/kube-apiserver/app/apiextensions.go | APIExtensions 配置与创建 |
cmd/kube-apiserver/app/options/ | 命令行选项定义与验证 |
pkg/kubeapiserver/ | kube-apiserver 特有配置(认证/授权/存储/准入) |
pkg/controlplane/ | 核心控制面实例(Instance)、内置 API 注册 |
staging/k8s.io/apiserver/ | 通用 API Server 框架(GenericAPIServer) |
staging/k8s.io/apiextensions-apiserver/ | CRD API Server |
staging/k8s.io/kube-aggregator/ | API 聚合 Server |
2.2 核心类结构与接口定义
2.2.1 顶层入口结构
// cmd/kube-apiserver/apiserver.go — 进程入口
func main() {
rand.Seed(time.Now().UnixNano())
command := app.NewAPIServerCommand() // 构建 cobra.Command
logs.InitLogs()
defer logs.FlushLogs()
command.Execute() // 执行命令
}
NewAPIServerCommand() 构建 cobra 命令,其 RunE 函数完成:
Complete(s)— 填充默认值,生成completedServerRunOptionsValidate()— 验证选项合法性Run(completedOptions, stopCh)— 创建并运行服务器
2.2.2 ServerRunOptions(选项聚合根)
type ServerRunOptions struct {
GenericServerRunOptions *genericoptions.ServerRunOptions // 通用服务选项
Etcd *genericoptions.EtcdOptions // etcd 存储选项
SecureServing *genericoptions.SecureServingOptionsWithLoopback
Audit *genericoptions.AuditOptions
Features *genericoptions.FeatureOptions
Admission *kubeoptions.AdmissionOptions // 准入控制选项
Authentication *kubeoptions.BuiltInAuthenticationOptions // 认证选项
Authorization *kubeoptions.BuiltInAuthorizationOptions // 授权选项
CloudProvider *kubeoptions.CloudProviderOptions
APIEnablement *genericoptions.APIEnablementOptions
EgressSelector *genericoptions.EgressSelectorOptions
// ... 业务特定字段
ServiceClusterIPRanges string
ServiceNodePortRange utilnet.PortRange
MasterCount int
// ... ServiceAccount 相关
}
2.2.3 三层服务器结构
GenericAPIServer — 通用 API 服务器基类:
type GenericAPIServer struct {
discoveryAddresses discovery.Addresses
LoopbackClientConfig *restclient.Config
SecureServingInfo *SecureServingInfo
Handler *APIServerHandler // 请求处理器链
admissionControl admission.Interface
Authorizer authorizer.Authorizer
AuditBackend audit.Backend
delegationTarget DelegationTarget // 委托目标
postStartHooks map[string]postStartHookEntry
DiscoveryGroupManager discovery.GroupManager
// ...
}
DelegationTarget — 委托接口:
type DelegationTarget interface {
UnprotectedHandler() http.Handler
PostStartHooks() map[string]postStartHookEntry
PreShutdownHooks() map[string]preShutdownHookEntry
HealthzChecks() []healthz.HealthChecker
ListedPaths() []string
NextDelegate() DelegationTarget
PrepareRun() preparedGenericAPIServer
}
Instance(控制面实例):
type Instance struct {
GenericAPIServer *genericapiserver.GenericAPIServer
ClusterAuthenticationInfo clusterauthenticationtrust.ClusterAuthenticationInfo
}
APIAggregator:
type APIAggregator struct {
GenericAPIServer *genericapiserver.GenericAPIServer
delegateHandler http.Handler
proxyHandlers map[string]*proxyHandler
handledGroups sets.String
lister listers.APIServiceLister
APIRegistrationInformers informers.SharedInformerFactory
serviceResolver ServiceResolver
// ...
}
CustomResourceDefinitions:
type CustomResourceDefinitions struct {
GenericAPIServer *genericapiserver.GenericAPIServer
Informers externalinformers.SharedInformerFactory
}
2.3 核心方法清单
| 方法 | 位置 | 作用 |
|---|---|---|
NewAPIServerCommand() | app/server.go | 构造 cobra 命令 |
Run() | app/server.go | 服务器运行入口 |
CreateServerChain() | app/server.go | 构建三层委托服务器链 |
CreateKubeAPIServerConfig() | app/server.go | 构建核心 API 服务器配置 |
buildGenericConfig() | app/server.go | 构建通用配置(认证/授权/存储/准入) |
BuildAuthorizer() | app/server.go | 构建授权器 |
createAPIExtensionsConfig() | app/apiextensions.go | 构建 CRD 服务器配置 |
createAPIExtensionsServer() | app/apiextensions.go | 创建 CRD 服务器实例 |
createAggregatorConfig() | app/aggregator.go | 构建聚合服务器配置 |
createAggregatorServer() | app/aggregator.go | 创建聚合服务器实例 |
Config.New() | genericapiserver/config.go | 创建 GenericAPIServer 实例 |
DefaultBuildHandlerChain() | genericapiserver/config.go | 构建请求处理过滤器链 |
PrepareRun() | genericapiserver/genericapiserver.go | 准备运行(安装健康检查/OpenAPI) |
Run() | genericapiserver/genericapiserver.go | 启动 HTTPS 服务 |
InstallLegacyAPI() | controlplane/instance.go | 安装核心 v1 API |
InstallAPIs() | controlplane/instance.go | 安装所有 API Group |
authenticator.Config.New() | kubeapiserver/authenticator/config.go | 构建认证器链 |
authorizer.Config.New() | kubeapiserver/authorizer/config.go | 构建授权器链 |
admission.Config.New() | kubeapiserver/admission/config.go | 构建准入插件初始化器 |
2.4 内部调用关系
2.5 数据流入流出方式
数据流入:
- HTTPS 请求 → Filter Chain → 认证/授权/准入 → REST Handler → Registry → etcd
数据流出:
- etcd 响应 → Codec 反序列化 → REST Handler 序列化 → HTTP Response
- Watch 机制:etcd Watch → Cacher → WatchHandler → Chunked Transfer Encoding
三、核心业务逻辑深度解析
3.1 完整启动流程
3.1.1 入口点(apiserver.go)
func main() {
rand.Seed(time.Now().UnixNano()) // 随机种子初始化
pflag.CommandLine.SetNormalizeFunc(cliflag.WordSepNormalizeFunc) // 标准化 flag 名
command := app.NewAPIServerCommand() // 构建命令
logs.InitLogs() // 日志初始化
defer logs.FlushLogs()
command.Execute() // 执行
}
3.1.2 命令构造(server.go — NewAPIServerCommand)
func NewAPIServerCommand() *cobra.Command {
s := options.NewServerRunOptions() // 创建默认选项
cmd := &cobra.Command{
Use: "kube-apiserver",
RunE: func(cmd *cobra.Command, args []string) error {
// 1. 打印版本
verflag.PrintAndExitIfRequested()
// 2. 打印所有 flag 值
cliflag.PrintFlags(fs)
// 3. 检查 insecure-port 必须为 0
checkNonZeroInsecurePort(fs)
// 4. Complete — 填充默认值
completedOptions, err := Complete(s)
// 5. Validate — 验证配置
completedOptions.Validate()
// 6. Run — 启动服务
return Run(completedOptions, genericapiserver.SetupSignalHandler())
},
}
// 注册所有 flag set
namedFlagSets := s.Flags()
for _, f := range namedFlagSets.FlagSets {
fs.AddFlagSet(f)
}
return cmd
}
3.1.3 Complete 阶段(server.go — Complete)
Complete 阶段关键逻辑:
- 默认广告地址:
s.GenericServerRunOptions.DefaultAdvertiseAddress() - Service ClusterIP 解析:解析
--service-cluster-ip-range为 Primary/Secondary 范围 - 自签名证书:
s.SecureServing.MaybeDefaultWithSelfSignedCerts()— 若未提供证书则生成自签名 - ExternalHost 推断:从 AdvertiseAddress 或 hostname 推断
- ServiceAccount 密钥处理:
- 若未指定
--service-account-signing-key-file,默认使用 TLS 私钥 - 构造
JWTTokenGenerator用于签发 ServiceAccount Token
- 若未指定
- WatchCache 大小配置:解析
--watch-cache-sizes覆盖默认值 - RuntimeConfig 规范化:将
v1、api/v1统一映射为/v1
3.1.4 Run 阶段
func Run(completeOptions completedServerRunOptions, stopCh <-chan struct{}) error {
klog.Infof("Version: %+v", version.Get()) // 打印版本信息
server, err := CreateServerChain(completeOptions, stopCh) // 构建服务器链
prepared, err := server.PrepareRun() // 准备运行
return prepared.Run(stopCh) // 启动 HTTPS
}
3.1.5 CreateServerChain — 三层委托链构建
这是整个启动流程中最核心的编排函数:
func CreateServerChain(completedOptions completedServerRunOptions, stopCh <-chan struct{}) (*aggregatorapiserver.APIAggregator, error) {
// Step 1: 创建 NodeDialer(SSH 隧道支持,已废弃)
nodeTunneler, proxyTransport, err := CreateNodeDialer(completedOptions)
// Step 2: 构建 KubeAPIServer 配置
kubeAPIServerConfig, serviceResolver, pluginInitializer, err :=
CreateKubeAPIServerConfig(completedOptions, nodeTunneler, proxyTransport)
// Step 3: 构建 APIExtensions (CRD) 配置和服务器
apiExtensionsConfig, err := createAPIExtensionsConfig(
*kubeAPIServerConfig.GenericConfig, ...)
apiExtensionsServer, err := createAPIExtensionsServer(
apiExtensionsConfig, genericapiserver.NewEmptyDelegate())
// CRD → 委托给 EmptyDelegate(最底层)
// Step 4: 创建 KubeAPIServer(核心内置 API)
kubeAPIServer, err := CreateKubeAPIServer(
kubeAPIServerConfig, apiExtensionsServer.GenericAPIServer)
// KubeAPIServer → 委托给 CRD Server
// Step 5: 构建 Aggregator 配置和服务器
aggregatorConfig, err := createAggregatorConfig(
*kubeAPIServerConfig.GenericConfig, ...)
aggregatorServer, err := createAggregatorServer(
aggregatorConfig, kubeAPIServer.GenericAPIServer, apiExtensionsServer.Informers)
// Aggregator → 委托给 KubeAPIServer
return aggregatorServer, nil
}
委托链方向(请求处理顺序):
Aggregator → KubeAPIServer → APIExtensions(CRD) → EmptyDelegate(404)
3.2 请求处理链路
3.2.1 Handler Chain 构建
DefaultBuildHandlerChain 是请求处理的核心过滤器链,在 config.go 中定义:
func DefaultBuildHandlerChain(apiHandler http.Handler, c *Config) http.Handler {
// 从内到外逐层包装(最后包装的最先执行)
// 1. Authorization — 授权检查
handler = genericapifilters.WithAuthorization(handler, c.Authorization.Authorizer, c.Serializer)
// 2. PriorityAndFairness 或 MaxInFlight — 限流
if c.FlowControl != nil {
handler = genericfilters.WithPriorityAndFairness(handler, c.LongRunningFunc, c.FlowControl)
} else {
handler = genericfilters.WithMaxInFlightLimit(handler, ...)
}
// 3. Impersonation — 身份伪装检查
handler = genericapifilters.WithImpersonation(handler, c.Authorization.Authorizer, c.Serializer)
// 4. Audit — 审计日志
handler = genericapifilters.WithAudit(handler, c.AuditBackend, c.AuditPolicyChecker, c.LongRunningFunc)
// 5. Authentication — 认证
failedHandler := genericapifilters.Unauthorized(c.Serializer)
handler = genericapifilters.WithAuthentication(handler, c.Authentication.Authenticator, failedHandler, c.Authentication.APIAudiences)
// 6. RequestInfo — 解析请求信息
handler = genericapifilters.WithRequestInfo(handler, c.RequestInfoResolver)
// 7. RequestDeadline — 请求超时
handler = genericapifilters.WithRequestDeadline(handler, ...)
// 8. WaitGroup — 请求追踪
handler = genericfilters.WithWaitGroup(handler, c.LongRunningFunc, c.HandlerChainWaitGroup)
// 9. CacheControl / HSTS / PanicRecovery 等
handler = genericapifilters.WithCacheControl(handler)
handler = genericfilters.WithHSTS(handler, c.HSTSDirectives)
handler = genericapifilters.WithRequestReceivedTimestamp(handler)
handler = genericfilters.WithPanicRecovery(handler, c.RequestInfoResolver)
return handler
}
执行顺序(外到内,即请求最先经过的):
PanicRecovery → RequestReceivedTimestamp → HSTS → CacheControl →
WaitGroup → RequestDeadline → RequestInfo → Authentication →
Audit → Impersonation → PriorityAndFairness → Authorization → apiHandler
3.2.2 Director 路由分发
请求经过 Handler Chain 后进入 APIServerHandler.Director:
func (d director) ServeHTTP(w http.ResponseWriter, req *http.Request) {
path := req.URL.Path
// 检查 go-restful 注册的 WebService 是否匹配
for _, ws := range d.goRestfulContainer.RegisteredWebServices() {
if ws.RootPath() == "/apis" {
// /apis 和 /apis/ 需特殊处理(discovery)
if path == "/apis" || path == "/apis/" {
d.goRestfulContainer.Dispatch(w, req)
return
}
}
if strings.HasPrefix(path, ws.RootPath()) {
// 精确匹配或路径边界匹配
d.goRestfulContainer.Dispatch(w, req)
return
}
}
// 未匹配 go-restful,走 nonGoRestfulMux(包含委托处理)
d.nonGoRestfulMux.ServeHTTP(w, req)
}
在 Aggregator 层,nonGoRestfulMux 上的 /apis/ 前缀由 proxyHandler 处理,负责将请求代理到后端扩展 API Server。
3.3 API Group/Version 路由机制
安装过程
在 controlplane/instance.go 中:
- InstallLegacyAPI — 安装
/api/v1核心组(Pods、Services、Nodes 等) - InstallAPIs — 安装所有 API Group(apps、batch、rbac 等)
每个 RESTStorageProvider 提供 NewRESTStorage() 方法,返回 APIGroupInfo:
type APIGroupInfo struct {
PrioritizedVersions []schema.GroupVersion
VersionedResourcesStorageMap map[string]map[string]rest.Storage // version → resource → storage
Scheme *runtime.Scheme
NegotiatedSerializer runtime.NegotiatedSerializer
// ...
}
InstallAPIs 遍历所有 Provider,调用 GenericAPIServer.InstallAPIGroup() 注册路由。
3.4 etcd 存储层架构
3.4.1 StorageFactory 构建
在 buildGenericConfig 中:
storageFactoryConfig := kubeapiserver.NewStorageFactoryConfig()
storageFactoryConfig.APIResourceConfig = genericConfig.MergedResourceConfig
completedStorageFactoryConfig, err := storageFactoryConfig.Complete(s.Etcd)
storageFactory, lastErr = completedStorageFactoryConfig.New()
NewStorageFactoryConfig() 设置了特殊资源前缀映射(SpecialDefaultResourcePrefixes):
var SpecialDefaultResourcePrefixes = map[schema.GroupResource]string{
{Group: "", Resource: "replicationcontrollers"}: "controllers",
{Group: "", Resource: "endpoints"}: "services/endpoints",
{Group: "", Resource: "nodes"}: "minions",
{Group: "", Resource: "services"}: "services/specs",
{Group: "networking.k8s.io", Resource: "ingresses"}: "ingress",
// ...
}
这是历史兼容性映射——etcd 中的 key 前缀与 API Group/Resource 名不一定相同。
3.4.2 存储层调用链
3.4.3 关键存储配置
- 默认存储媒体类型:
application/vnd.kubernetes.protobuf(二进制,高效) - WatchCache:默认启用,Events 资源除外(
DefaultWatchCacheSizes中 events 大小为 0) - 加密:支持
--encryption-provider-config对特定资源进行加密存储 - 分页:受
APIListChunkingfeature gate 控制 - Cohabitating Resources:多个 Group 的同名资源共享 etcd key 前缀
3.5 认证(Authentication)深度解析
3.5.1 配置构建
在 buildGenericConfig 中,认证配置通过 s.Authentication.ApplyTo() 应用:
s.Authentication.ApplyTo(&genericConfig.Authentication, genericConfig.SecureServing,
genericConfig.EgressSelector, genericConfig.OpenAPIConfig, clientgoExternalClient, versionedInformers)
该方法内部:
- 将
BuiltInAuthenticationOptions转换为kubeauthenticator.Config - 设置
ServiceAccountTokenGetter(从 Informer 获取 SA/Secret) - 设置
BootstrapTokenAuthenticator - 配置 EgressSelector 的 CustomDial
- 调用
authenticatorConfig.New()构建认证器
3.5.2 认证器链构建(authenticator/config.go — New)
func (config Config) New() (authenticator.Request, *spec.SecurityDefinitions, error) {
var authenticators []authenticator.Request
var tokenAuthenticators []authenticator.Token
// 1. Front-Proxy (RequestHeader) 认证
if config.RequestHeaderConfig != nil {
requestHeaderAuthenticator := headerrequest.NewDynamicVerifyOptionsSecure(...)
authenticators = append(authenticators, requestHeaderAuthenticator)
}
// 2. X509 客户端证书认证
if config.ClientCAContentProvider != nil {
certAuth := x509.NewDynamic(config.ClientCAContentProvider.VerifyOptions, x509.CommonNameUserConversion)
authenticators = append(authenticators, certAuth)
}
// 3. Token 文件认证
if len(config.TokenAuthFile) > 0 {
tokenAuth, _ := newAuthenticatorFromTokenFile(config.TokenAuthFile)
tokenAuthenticators = append(tokenAuthenticators, tokenAuth)
}
// 4. ServiceAccount (Legacy) 认证
if len(config.ServiceAccountKeyFiles) > 0 {
serviceAccountAuth, _ := newLegacyServiceAccountAuthenticator(...)
tokenAuthenticators = append(tokenAuthenticators, serviceAccountAuth)
}
// 5. ServiceAccount (JWT) 认证
if config.ServiceAccountIssuer != "" {
serviceAccountAuth, _ := newServiceAccountAuthenticator(...)
tokenAuthenticators = append(tokenAuthenticators, serviceAccountAuth)
}
// 6. Bootstrap Token 认证
if config.BootstrapToken {
tokenAuthenticators = append(tokenAuthenticators, config.BootstrapTokenAuthenticator)
}
// 7. OIDC 认证
if len(config.OIDCIssuerURL) > 0 && len(config.OIDCClientID) > 0 {
oidcAuth, _ := newAuthenticatorFromOIDCIssuerURL(oidc.Options{...})
tokenAuthenticators = append(tokenAuthenticators, oidcAuth)
}
// 8. Webhook Token 认证
if len(config.WebhookTokenAuthnConfigFile) > 0 {
webhookTokenAuth, _ := newWebhookTokenAuthenticator(config)
tokenAuthenticators = append(tokenAuthenticators, webhookTokenAuth)
}
// 合并 Token 认证器
if len(tokenAuthenticators) > 0 {
tokenAuth := tokenunion.New(tokenAuthenticators...)
// 可选缓存
if config.TokenSuccessCacheTTL > 0 || config.TokenFailureCacheTTL > 0 {
tokenAuth = tokencache.New(tokenAuth, true, ...)
}
authenticators = append(authenticators,
bearertoken.New(tokenAuth),
websocket.NewProtocolAuthenticator(tokenAuth))
}
// 合并所有认证器
authenticator := union.New(authenticators...)
// 自动添加 system:authenticated 组
authenticator = group.NewAuthenticatedGroupAdder(authenticator)
// 匿名认证兜底
if config.Anonymous {
authenticator = union.NewFailOnError(authenticator, anonymous.NewAuthenticator())
}
return authenticator, &securityDefinitions, nil
}
3.6 授权(Authorization)深度解析
3.6.1 配置构建
func BuildAuthorizer(s *options.ServerRunOptions, EgressSelector *egressselector.EgressSelector,
versionedInformers clientgoinformers.SharedInformerFactory) (authorizer.Authorizer, authorizer.RuleResolver, error) {
authorizationConfig := s.Authorization.ToAuthorizationConfig(versionedInformers)
if EgressSelector != nil {
egressDialer, _ := EgressSelector.Lookup(egressselector.ControlPlane.AsNetworkContext())
authorizationConfig.CustomDial = egressDialer
}
return authorizationConfig.New()
}
3.6.2 授权器链构建(authorizer/config.go — New)
func (config Config) New() (authorizer.Authorizer, authorizer.RuleResolver, error) {
// 必须至少指定一个授权模式
if len(config.AuthorizationModes) == 0 {
return nil, nil, fmt.Errorf("at least one authorization mode must be passed")
}
for _, authorizationMode := range config.AuthorizationModes {
switch authorizationMode {
case modes.ModeNode:
// Node 授权器 — 专门授权 kubelet 请求
graph := node.NewGraph()
node.AddGraphEventHandlers(graph, ...) // 注册 Node/Pod/PV/VA Informer
nodeAuthorizer := node.NewAuthorizer(graph, nodeidentifier.NewDefaultNodeIdentifier(), bootstrappolicy.NodeRules())
case modes.ModeRBAC:
// RBAC 授权器 — 基于角色绑定
rbacAuthorizer := rbac.New(
&rbac.RoleGetter{Lister: ...},
&rbac.RoleBindingLister{Lister: ...},
&rbac.ClusterRoleGetter{Lister: ...},
&rbac.ClusterRoleBindingLister{Lister: ...},
)
case modes.ModeABAC:
// ABAC 授权器 — 基于策略文件
abacAuthorizer, _ := abac.NewFromFile(config.PolicyFile)
case modes.ModeWebhook:
// Webhook 授权器 — 远程服务决策
webhookAuthorizer, _ := webhook.New(config.WebhookConfigFile, ...)
case modes.ModeAlwaysAllow:
// AlwaysAllow — 全部放行(仅用于测试)
case modes.ModeAlwaysDeny:
// AlwaysDeny — 全部拒绝(仅用于测试)
}
}
// 返回 Union 授权器 — 按顺序检查,任一通过即放行
return union.New(authorizers...), union.NewRuleResolvers(ruleResolvers...), nil
}
推荐授权模式:--authorization-mode=Node,RBAC
3.7 Admission Controller 链路深度解析
3.7.1 插件注册
在 options/plugins.go 中,RegisterAllAdmissionPlugins 注册了所有内置插件,AllOrderedPlugins 定义了执行顺序:
var AllOrderedPlugins = []string{
admit.PluginName, // AlwaysAdmit
autoprovision.PluginName, // NamespaceAutoProvision
lifecycle.PluginName, // NamespaceLifecycle
exists.PluginName, // NamespaceExists
scdeny.PluginName, // SecurityContextDeny
antiaffinity.PluginName, // LimitPodHardAntiAffinityTopology
limitranger.PluginName, // LimitRanger
serviceaccount.PluginName, // ServiceAccount
noderestriction.PluginName, // NodeRestriction
nodetaint.PluginName, // TaintNodesByCondition
alwayspullimages.PluginName, // AlwaysPullImages
// ... 更多内置插件 ...
mutatingwebhook.PluginName, // MutatingAdmissionWebhook
validatingwebhook.PluginName, // ValidatingAdmissionWebhook
resourcequota.PluginName, // ResourceQuota
deny.PluginName, // AlwaysDeny
}
默认启用的插件(DefaultOffAdmissionPlugins 的补集):
- NamespaceLifecycle、LimitRanger、ServiceAccount、DefaultStorageClass
- PersistentVolumeClaimResize、DefaultTolerationSeconds
- MutatingAdmissionWebhook、ValidatingAdmissionWebhook、ResourceQuota
- StorageObjectInUseProtection、PodPriority、TaintNodesByCondition
- RuntimeClass、CertificateApproval/Signing/SubjectRestriction、DefaultIngressClass
3.7.2 插件初始化
// kubeapiserver/admission/config.go
func (c *Config) New(proxyTransport, egressSelector, serviceResolver) ([]admission.PluginInitializer, PostStartHookFunc, error) {
// Webhook 初始化器 — 提供 auth 解析和 service 解析
webhookPluginInitializer := webhookinit.NewPluginInitializer(webhookAuthResolverWrapper, serviceResolver)
// Kubernetes 特有初始化器
kubePluginInitializer := NewPluginInitializer(cloudConfig, discoveryRESTMapper, quotaConfiguration)
// PostStartHook — 定期重置 Discovery REST Mapper
admissionPostStartHook := func(context PostStartHookContext) error {
discoveryRESTMapper.Reset()
go utilwait.Until(discoveryRESTMapper.Reset, 30*time.Second, context.StopCh)
return nil
}
return []admission.PluginInitializer{webhookPluginInitializer, kubePluginInitializer}, admissionPostStartHook, nil
}
PluginInitializer.Initialize() 通过接口注入:
func (i *PluginInitializer) Initialize(plugin admission.Interface) {
if wants, ok := plugin.(WantsCloudConfig); ok {
wants.SetCloudConfig(i.cloudConfig)
}
if wants, ok := plugin.(WantsRESTMapper); ok {
wants.SetRESTMapper(i.restMapper)
}
if wants, ok := plugin.(initializer.WantsQuotaConfiguration); ok {
wants.SetQuotaConfiguration(i.quotaConfiguration)
}
}
3.7.3 Admission 执行流程
Admission 分两阶段执行:
- Mutating 阶段:所有 Mutating 插件按顺序执行,可修改对象
- Validating 阶段:所有 Validating 插件按顺序执行,只读验证
3.8 CRD 处理流程
3.8.1 APIExtensions Server 创建
在 createAPIExtensionsConfig 中:
- 浅拷贝 GenericConfig,清空 PostStartHooks 和 RESTOptionsGetter
- 设置 CRD 专用的 Codec(
apiextensionsapiserver.Codecs.LegacyCodec) - 设置 EncodeVersioner(优先 v1beta1 编码以节省存储空间)
- 配置
CRDRESTOptionsGetter(CRD 资源自身的 etcd 存储) - 配置 ServiceResolver(用于 Webhook 转换)
3.8.2 CRD Handler 注册
在 apiextensions-apiserver 的 New() 方法中:
crdHandler, err := NewCustomResourceDefinitionHandler(
versionDiscoveryHandler,
groupDiscoveryHandler,
s.Informers.Apiextensions().V1().CustomResourceDefinitions(),
delegateHandler,
c.ExtraConfig.CRDRESTOptionsGetter,
c.GenericConfig.AdmissionControl,
establishingController,
c.ExtraConfig.ServiceResolver,
c.ExtraConfig.AuthResolverWrapper,
c.ExtraConfig.MasterCount,
s.GenericAPIServer.Authorizer,
// ...
)
// CRD Handler 拦截 /apis 下的所有请求
s.GenericAPIServer.Handler.NonGoRestfulMux.Handle("/apis", crdHandler)
s.GenericAPIServer.Handler.NonGoRestfulMux.HandlePrefix("/apis/", crdHandler)
这意味着 CRD Handler 在 NonGoRestfulMux 中优先于 go-restful 处理 /apis/ 下的请求。
3.9 API 聚合(Aggregation)架构
3.9.1 Aggregator Server 创建
func createAggregatorServer(aggregatorConfig, delegateAPIServer, apiExtensionInformers) (*APIAggregator, error) {
aggregatorServer, err := aggregatorConfig.Complete().NewWithDelegate(delegateAPIServer)
// 创建自动注册控制器
autoRegistrationController := autoregister.NewAutoRegisterController(...)
apiServices := apiServicesToRegister(delegateAPIServer, autoRegistrationController)
// CRD 注册控制器 — 将 CRD 对应的 APIService 自动注册
crdRegistrationController := crdregistration.NewCRDRegistrationController(
apiExtensionInformers.Apiextensions().V1().CustomResourceDefinitions(),
autoRegistrationController)
// PostStartHook — 启动控制器
aggregatorServer.GenericAPIServer.AddPostStartHook("kube-apiserver-autoregistration", func(context) error {
go crdRegistrationController.Run(5, context.StopCh)
// 等待 CRD 初始同步完成后再启动 autoRegistration
go func() {
if aggregatorConfig.GenericConfig.MergedResourceConfig.AnyVersionForGroupEnabled("apiextensions.k8s.io") {
crdRegistrationController.WaitForInitialSync()
}
autoRegistrationController.Run(5, context.StopCh)
}()
return nil
})
}
3.9.2 APIService 路由
Aggregator 维护 proxyHandlers 映射,每个 APIService 对应一个 proxyHandler:
- Local APIService(如 v1. 核心组)→ 委托给内置 KubeAPIServer
- External APIService(如 metrics.k8s.io/v1beta1)→ 通过 ServiceResolver 解析后端,代理请求
3.9.3 APIService 优先级
在 aggregator.go 中定义了 apiVersionPriorities,控制 Discovery 输出中 API Group 的排序:
var apiVersionPriorities = map[schema.GroupVersion]priority{
{Group: "", Version: "v1"}: {group: 18000, version: 1},
{Group: "apps", Version: "v1"}: {group: 17800, version: 15},
{Group: "rbac.authorization.k8s.io", Version: "v1"}: {group: 17000, version: 15},
// ... group 值越大优先级越高
}
3.10 Watch 机制流程
Watch 机制的关键组件:
-
Cacher — 位于
etcd3.store和 etcd client 之间,维护内存缓存- 所有 Watch 请求由 Cacher 服务,而非直接打 etcd
- Cacher 内部有一个 reflector 从 etcd List/Watch 同步数据
- Events 资源默认禁用 WatchCache(
DefaultWatchCacheSizes中 events 大小为 0)
-
WatchCache — 环形缓冲区,存储最近的变更事件
- 默认容量 100(可通过
--watch-cache-sizes调整) - 新 Watcher 从缓存中重放历史,避免 etcd 压力
- 默认容量 100(可通过
-
Bookmark — 定期发送
BOOKMARK类型事件,包含当前 ResourceVersion- 帮助客户端维持 Watch 连接的 ResourceVersion
-
超时机制 —
MinRequestTimeout(默认 1800 秒)控制 Watch 的随机超时
四、Mermaid 图汇总
图1:kube-apiserver 启动流程图
图2:请求处理链路图
图3:API Group/Version 路由图
图4:etcd 存储层架构图
图5:Admission Controller 链路图
图6:认证(Authentication)流程图
图7:授权(Authorization)流程图
图8:CRD 处理流程图
图9:API 聚合(Aggregation)架构图
图10:Watch 机制流程图
五、关键设计模式与深度剖析
5.1 委托模式(Delegation Pattern)
kube-apiserver 的三层结构采用委托模式:
- 每层服务器持有
delegationTarget引用 - 请求在本层无法处理时,委托给下一层
- 通过
UnprotectedHandler()获取下一层的裸 Handler(无 Filter Chain 包装) - Aggregator 层在最外层,收到所有请求后决定路由
这个设计的关键意义在于:每个 API Server 独立拥有自己的 Filter Chain(认证/授权/准入),但共享同一个 HTTPS 监听端口。
5.2 Handler Chain 构建顺序
DefaultBuildHandlerChain 采用洋葱模型:
最外层 → PanicRecovery
→ RequestTimestamp
→ WaitGroup
→ RequestDeadline
→ RequestInfo
→ Authentication ← 最先认证
→ Audit
→ Impersonation
→ PriorityAndFairness
→ Authorization ← 认证后授权
最内层 → apiHandler
关键设计决策:
- Authentication 在 Authorization 之前:先确认身份,再判断权限
- Audit 包裹 Auth 和 Impersonation:记录认证和伪装结果
- PriorityAndFairness 在 Authorization 之前:限流优先,避免授权计算消耗资源
- Impersonation 在 Authorization 之前:伪装后的身份用于授权决策
5.3 认证器的 Union 模式
认证器采用短路模式:Union 中的认证器按顺序尝试,第一个成功即返回。但 Anonymous 作为兜底,通过 NewFailOnError 包装——只有前序认证器明确拒绝时才走匿名认证。
// 关键逻辑:
authenticator = union.New(authenticators...) // 任一成功即通过
authenticator = group.NewAuthenticatedGroupAdder(authenticator) // 自动加组
if config.Anonymous {
authenticator = union.NewFailOnError(authenticator, anonymous.NewAuthenticator())
// NewFailOnError: 前序认证器返回错误时不走匿名,仅前序返回false时走匿名
}
5.4 授权器的 Union 模式
授权器同样采用 Union 模式,但语义为任一通过即放行(OR 语义):
return union.New(authorizers...), union.NewRuleResolvers(ruleResolvers...), nil
这意味着 --authorization-mode=Node,RBAC 等价于"Node授权器允许 OR RBAC授权器允许"。
5.5 StorageFactory 的资源映射
SpecialDefaultResourcePrefixes 保留了 etcd 历史兼容性。例如:
services→services/specs(而非services)nodes→minions(历史遗留名称)replicationcontrollers→controllers
AddCohabitatingResources 将跨 Group 的同名资源映射到同一 etcd 前缀:
storageFactory.AddCohabitatingResources(apps.Resource("deployments"), extensions.Resource("deployments"))
这确保 apps/v1/deployments 和 extensions/v1beta1/deployments 操作同一个 etcd 数据。
5.6 ServiceAccount Token 签发
Complete 阶段的 ServiceAccount Token 签发逻辑:
if s.ServiceAccountSigningKeyFile != "" && s.Authentication.ServiceAccounts.Issuer != "" {
sk, _ := keyutil.PrivateKeyFromFile(s.ServiceAccountSigningKeyFile)
s.ServiceAccountIssuer, _ = serviceaccount.JWTTokenGenerator(
s.Authentication.ServiceAccounts.Issuer, sk)
s.ServiceAccountTokenMaxExpiration = s.Authentication.ServiceAccounts.MaxExpiration
}
关键约束:
MaxExpiration必须在[1h, 2^32s]范围内ExtendExpiration启用时,注入的 Token 最长可延至 1 年,以平滑过渡
5.7 API 优先级与公平性(FlowControl)
当 APIPriorityAndFairness feature gate 启用且 --enable-priority-and-fairness=true 时:
genericConfig.FlowControl = BuildPriorityAndFairness(s, clientgoExternalClient, versionedInformers)
// 实现:
return utilflowcontrol.New(
versionedInformer,
extclient.FlowcontrolV1beta1(),
s.GenericServerRunOptions.MaxRequestsInFlight + s.GenericServerRunOptions.MaxMutatingRequestsInFlight,
s.GenericServerRunOptions.RequestTimeout/4,
)
FlowControl 替代了传统的 MaxRequestsInFlight 限流,提供更细粒度的优先级队列。
5.8 优雅关闭(Graceful Shutdown)
preparedGenericAPIServer.Run() 实现了优雅关闭:
func (s preparedGenericAPIServer) Run(stopCh <-chan struct{}) error {
delayedStopCh := make(chan struct{})
go func() {
<-stopCh
close(s.readinessStopCh) // /readyz 开始返回失败
time.Sleep(s.ShutdownDelayDuration) // 等待LB感知
close(delayedStopCh) // 关闭监听器
}()
stoppedCh, _ := s.NonBlockingRun(delayedStopCh) // 启动HTTPS
<-stopCh // 等待信号
s.RunPreShutdownHooks() // 执行预关闭钩子
<-delayedStopCh // 等待延迟关闭
<-stoppedCh // 等待连接完成
s.HandlerChainWaitGroup.Wait() // 等待所有请求完成
}
流程:stopCh → readinessStopCh(readyz失败) → ShutdownDelayDuration → 关闭Listener → PreShutdownHooks → 等待连接结束
5.9 PostStartHook 机制
kube-apiserver 在启动后通过 PostStartHook 异步启动多种控制器和初始化逻辑:
| Hook 名称 | 作用 |
|---|---|
generic-apiserver-start-informers | 启动 SharedInformerFactory |
start-kube-apiserver-admission-initializer | 准入控制器初始化 |
start-kube-aggregator-informers | Aggregator Informer 启动 |
apiservice-registration-controller | APIService 注册控制器 |
apiservice-status-available-controller | APIService 可用性检查 |
kube-apiserver-autoregistration | CRD/内置 API 自动注册 |
start-apiextensions-informers | CRD Informer 启动 |
start-apiextensions-controllers | CRD 控制器启动 |
crd-informer-synced | CRD Informer 同步等待 |
start-cluster-authentication-info-controller | 集群认证信息控制器 |
start-kube-apiserver-identity-lease-controller | APIServer 身份 Lease |
priority-and-fairness-config-consumer | FlowControl 启动 |
aggregator-reload-proxy-client-cert | 代理证书热更新 |
5.10 Loopback Client 配置
kube-apiserver 创建一个自环客户端(Loopback Client),用于:
- Admission Webhook 调用自身
- Informer List/Watch 自身资源
- PostStartHook 中操作自身 API
关键配置:
genericConfig.LoopbackClientConfig.ContentConfig.ContentType = "application/vnd.kubernetes.protobuf"
genericConfig.LoopbackClientConfig.DisableCompression = true
AuthorizeClientBearerToken 确保 loopback token 获得最高权限:
func AuthorizeClientBearerToken(loopback *restclient.Config, authn *AuthenticationInfo, authz *AuthorizationInfo) {
// 将 loopback token 注册为特权 token
// 通过 authenticatorfactory 和 authorizerfactory 的特例处理
}
六、配置选项与 Feature Gate
6.1 关键命令行参数
| 参数 | 默认值 | 说明 |
|---|---|---|
--etcd-servers | 无(必填) | etcd 服务器列表 |
--secure-port | 6443 | HTTPS 监听端口 |
--authorization-mode | AlwaysAllow | 授权模式(生产建议 Node,RBAC) |
--enable-admission-plugins | 见默认列表 | 启用的准入插件 |
--service-cluster-ip-range | 10.0.0.0/24 | Service ClusterIP 范围 |
--service-account-issuer | 无(必填) | SA Token 签发者 URL |
--service-account-key-file | TLS key | SA Token 验证公钥 |
--service-account-signing-key-file | 无 | SA Token 签名私钥 |
--client-ca-file | 无 | 客户端证书 CA |
--requestheader-client-ca-file | 无 | Front-Proxy CA |
--enable-aggregator-routing | false | Aggregator 路由模式 |
--watch-cache | true | 启用 WatchCache |
--max-requests-inflight | 400 | 最大并发只读请求 |
--max-mutating-requests-inflight | 200 | 最大并发变更请求 |
6.2 Feature Gate 影响
| Feature Gate | 影响 |
|---|---|
APIPriorityAndFairness | 启用 FlowControl 替代 MaxInFlight |
APIServerIdentity | APIServer 身份 Lease |
StorageVersionAPI | 存储版本迁移 |
APIListChunking | 分页 List |
ServiceAccountIssuerDiscovery | OIDC Discovery 端点 |
EgressSelector | 出站流量控制 |
七、总结
kube-apiserver 的架构核心在于三层委托链 + Handler Chain 洋葱模型:
- Aggregator 作为最外层网关,负责请求路由和扩展 API Server 代理
- KubeAPIServer 承载所有内置 API Group 的 REST Storage
- APIExtensions 提供 CRD 动态注册能力
请求处理链严格遵循认证→限流→授权→准入→验证→存储的顺序,每层都有清晰的职责边界。
存储层通过 StorageFactory 抽象 etcd 操作,Cacher 提供 Watch 缓存,Codec 优先使用 protobuf 编码,EncryptionConfig 支持静态加密。
认证支持 8+ 种机制(X509/TokenFile/SA/OIDC/Webhook/Bootstrap/RequestHeader/Anonymous),通过 Union 模式组合。
授权支持 6 种模式(Node/RBAC/ABAC/Webhook/AlwaysAllow/AlwaysDeny),通过 Union 模式组合。
准入控制分两阶段执行(Mutating→Validating),Webhook 扩展位于链末尾,ResourceQuota 作为最后防线。
这套架构使得 kube-apiserver 既是 Kubernetes 的 API 网关,也是其安全策略执行点和数据一致性保障核心。
更多推荐


所有评论(0)