eventRecorder
当集群中的node或pod异常时,大部分用户会使用kubectl查看对应的events,events本质上是k8s中的各个组件将运行时产生的各种事件汇报到apiserver,对于k8s中的可描述资源,使用kubectl describe都可以看到其相关的events。目前,k8s中包括controller-manager、kube-proxy、kube-scheduler、kubelet都使用了EventReconder,这里以kubelet的使用为例。
# 1.概述
# 1.1.定义
k8s的Event事件是一种资源对象,用于展示集群内发生的情况,k8s系统中的各个组件会将运行时发生的各种事件上报给apiserver。可以通过kubectl get event或kubectl describe命令查看资源对象的相关事件。apiserver会将Event事件存在etcd集群,为避免磁盘空间被填满,会强制执行暴露策略:最后一次事件发生后,删除1小时之前发生的事件。Events: Type Reason Age From Message ---- ------ ---- ---- ------- Normal Scheduled 19s default-scheduler Successfully assigned default/hpatest-bbb44c476-8d45v to 192.168.13.130 Normal Pulled 15s kubelet, 192.168.13.130 Container image "nginx" already present on machine Normal Created 15s kubelet, 192.168.13.130 Created container hpatest Normal Started 13s kubelet, 192.168.13.130 Started container hpatest1
2
3
4
5
6
7当集群中的
node或pod异常时,大部分用户会使用kubectl查看对应的events,对应以下事件打印程序。recorder.Eventf(cj, v1.EventTypeWarning, "FailedNeedsStart", "Cannot determine if job needs to be started: %v", err)1通过事件模块调用,可以确认基本上与
node或pod相关的模块或涉及到事件,如controller-manager、kube-proxy、kube-scheduler、kubelet等。
# 1.2.管理机制
Event事件的管理机制主要由三部分组成:1.
EventRecorder:事件生成者,k8s组件通过调用该方法生成事件2.
EventBroadcaster:事件广播器,负责消费EventRecorder产生的事件,然后分发给broadcasterWatcher3.
broadcasterWatcher:用于定义事件的处理方式,例如上报apiserver
# 2.原理分析
# 2.1.makeEventRecorder
kubelet初始化时会调用makeEventRecorder对Event进行Init操作。func makeEventRecorder(kubeDeps *kubelet.Dependencies, nodeName types.NodeName) { if kubeDeps.Recorder != nil { return } // 初始化 eventBroadcaster eventBroadcaster := record.NewBroadcaster() // 初始化 eventRecorder kubeDeps.Recorder = eventBroadcaster.NewRecorder(legacyscheme.Scheme, v1.EventSource{Component: componentKubelet, Host: string(nodeName)}) // 记录 event到日志 eventBroadcaster.StartStructuredLogging(3) if kubeDeps.EventClient != nil { klog.V(4).InfoS("Sending events to api server") // 上报 event 到apiserver eventBroadcaster.StartRecordingToSink(&v1core.EventSinkImpl{Interface: kubeDeps.EventClient.Events("")}) } else { klog.InfoS("No api server defined - no events will be sent to API server") } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18这个方法创建了一个
eventBroadcaster,这是一个事件广播器,会消费eventRecorder记录的事件并通过StartStructuredLogging和StartRecordingToSink分别将event发送给log和apiserver。
# 2.2.EventRecorder
# 2.2.1.事件记录
eventRecorder用作事件记录器,k8s系统组件通过它记录关键性事件。type EventRecorder interface { Event(object runtime.Object, eventtype, reason, message string) Eventf(object runtime.Object, eventtype, reason, messageFmt string, args ...interface{}) AnnotatedEventf(object runtime.Object, annotations map[string]string, eventtype, reason, messageFmt string, args ...interface{}) }1
2
3
4
5
6
7eventRecorder接口包含三个方法,Event用于记录刚发生的事件,Eventf通过使用fmt.Sprintf格式化输出事件的格式,AnnotatedEventf功能和Eventf一致,但是附加了注释字段。记录事件是一般是如下形式:recorder.Eventf(cj, v1.EventTypeWarning, "FailedNeedsStart", "Cannot determine if job needs to be started: %v", err)1Eventf会调用到EventRecorder的实现类recorderImpl,最后调用到generateEvent方法。
# 2.2.2.Event
func (recorder *recorderImpl) Event(object runtime.Object, eventtype, reason, message string) { recorder.generateEvent(object, nil, eventtype, reason, message) } func (recorder *recorderImpl) Eventf(object runtime.Object, eventtype, reason, messageFmt string, args ...interface{}) { recorder.Event(object, eventtype, reason, fmt.Sprintf(messageFmt, args...)) }1
2
3
4
5
6
7
# 2.2.3.generateEvent
func (recorder *recorderImpl) generateEvent(object runtime.Object, annotations map[string]string, eventtype, reason, message string) { ... // 实例化 event event := recorder.makeEvent(ref, annotations, eventtype, reason, message) event.Source = recorder.source // 调用ActionOrDrop方法,将事件写入incoming中 if sent := recorder.ActionOrDrop(watch.Added, event); !sent { klog.Errorf("unable to record event: too many queued events, dropped event %#v", event) } } func (m *Broadcaster) ActionOrDrop(action EventType, obj runtime.Object) bool { select { case m.incoming <- Event{action, obj}: return true default: return false } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
# 2.3.eventBroadcaster
# 2.3.1.事件广播
eventBroadcaster初始化的时候会调用NewBroadcaster方法。func NewBroadcaster() EventBroadcaster { return &eventBroadcasterImpl{ Broadcaster: watch.NewLongQueueBroadcaster(maxQueuedEvents, watch.DropIfChannelFull), sleepDuration: defaultSleepDuration, } }1
2
3
4
5
6这里会创建一个
eventBroadcasterImpl实例,并设置Broadcaster和sleepDuration两个字段,前者是这个方法的核心。func NewLongQueueBroadcaster(queueLength int, fullChannelBehavior FullChannelBehavior) *Broadcaster { m := &Broadcaster{ watchers: map[int64]*broadcasterWatcher{}, incoming: make(chan Event, queueLength), stopped: make(chan struct{}), watchQueueLength: queueLength, fullChannelBehavior: fullChannelBehavior, } m.distributing.Add(1) // 开启事件循环 go m.loop() return m }1
2
3
4
5
6
7
8
9
10
11
12
13这里初始化
broadcaster的时候,会初始化一个broadcasterWatcher,用于定义事件处理方式,例如上报apiserver等;初始化incoming管道,用于eventBroadcaster和eventRecorder进行事件传输。
# 2.3.2.loop
func (m *Broadcaster) loop() { // 获取m.incoming管道中的数据 for event := range m.incoming { if event.Type == internalRunFunctionMarker { event.Object.(functionFakeRuntimeObject)() continue } // 进行事件分发 m.distribute(event) } m.closeAll() m.distributing.Done() }1
2
3
4
5
6
7
8
9
10
11
12
13这个方法会一直在后台等待获取
m.incoming管道中的数据,然后调用distributing方法进行事件分发给broadcasterWatcher。incoming管道中的数据是eventRecorder调用Event方法传入的。
# 2.3.3.distribute
func (m *Broadcaster) distribute(event Event) { // 如果是非阻塞,那么使用DropIfChannelFull标识 if m.fullChannelBehavior == DropIfChannelFull { for _, w := range m.watchers { select { case w.result <- event: case <-w.stopped: default: // Don't block if the event can't be queued. } } } else { for _, w := range m.watchers { select { case w.result <- event: case <-w.stopped: } } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19如果是非阻塞,那么使用
DropIfChannelFull标识,在w.result管道满了以后,事件会丢失。如果没有default关键字,那么当w.result管道满了后,分发过程会阻塞并等待。这里之所以需要丢失事件,是因为随着k8s集群越来越大,上报事件也随之增多,那么每次上报都要对etcd进行读写,这样会给etcd集群带来读写压力。但是事件丢失并不会影响集群的正常工作,所以非阻塞分发机制下事件会丢失。
# 2.4.recordToSink
# 2.4.1.事件处理
makeEventRecorder调用StartRecordingToSink方法会将数据上报到apiserver。func (e *eventBroadcasterImpl) StartRecordingToSink(sink EventSink) watch.Interface { eventCorrelator := NewEventCorrelatorWithOptions(e.options) return e.StartEventWatcher( func(event *v1.Event) { recordToSink(sink, event, eventCorrelator, e.sleepDuration) }) } func (e *eventBroadcasterImpl) StartEventWatcher(eventHandler func(*v1.Event)) watch.Interface { watcher := e.Watch() go func() { defer utilruntime.HandleCrash() for watchEvent := range watcher.ResultChan() { event, ok := watchEvent.Object.(*v1.Event) if !ok { // This is all local, so there's no reason this should // ever happen. continue } // 回调传入的方法 eventHandler(event) } }() return watcher }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25StartRecordingToSink会调用StartEventWatcher方法,StartEventWatcher方法里面会异步的调用watcher.ResultChan()方法获取到broadcasterWatcher的result管道,result管道里的数据就是broadcaster的distribute进行分发的,最后会回调传入recordToSink方法。
# 2.4.2.recordToSink
func recordToSink(sink EventSink, event *v1.Event, eventCorrelator *EventCorrelator, sleepDuration time.Duration) { eventCopy := *event event = &eventCopy // 对事件作预处理,聚合相同事件 result, err := eventCorrelator.EventCorrelate(event) if err != nil { utilruntime.HandleError(err) } if result.Skip { return } tries := 0 for { // 把事件发送到apiserver if recordEvent(sink, result.Event, result.Patch, result.Event.Count > 1, eventCorrelator) { break } tries++ if tries >= maxTriesPerEvent { klog.Errorf("Unable to write event '%#v' (retry limit exceeded!)", event) break } if tries == 1 { time.Sleep(time.Duration(float64(sleepDuration) * rand.Float64())) } else { time.Sleep(sleepDuration) } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29recordToSink方法首先会调用EventCorrelate方法对event作预处理,聚合相同的事件,避免产生的事件过多,增加etcd和apiserver的压力,如果传入的event太多,那么result.Skip就会返回false。接下来会调用recordEvent方法把事件发送到apiserver,它会默认重试12次,并且每次重试都有一定时间间隔(默认10s)。
# 2.4.3.EventCorrelate
func (c *EventCorrelator) EventCorrelate(newEvent *v1.Event) (*EventCorrelateResult, error) { if newEvent == nil { return nil, fmt.Errorf("event is nil") } aggregateEvent, ckey := c.aggregator.EventAggregate(newEvent) observedEvent, patch, err := c.logger.eventObserve(aggregateEvent, ckey) if c.filterFunc(observedEvent) { return &EventCorrelateResult{Skip: true}, nil } return &EventCorrelateResult{Event: observedEvent, Patch: patch}, err }1
2
3
4
5
6
7
8
9
10
11EventCorrelate方法会调用EventAggregate、eventObserve进行聚合,调用filterFunc会调用到spamFilter.Filter方法进行过滤。func (e *EventAggregator) EventAggregate(newEvent *v1.Event) (*v1.Event, string) { now := metav1.NewTime(e.clock.Now()) var record aggregateRecord eventKey := getEventKey(newEvent) aggregateKey, localKey := e.keyFunc(newEvent) e.Lock() defer e.Unlock() // 查找缓存里边是否存在这样的记录 value, found := e.cache.Get(aggregateKey) if found { record = value.(aggregateRecord) } // maxIntervalInSeconds默认时间是600s,这里校验缓存里的记录是不是太老了,如果太老或找不到则会创建一个新的 maxInterval := time.Duration(e.maxIntervalInSeconds) * time.Second interval := now.Time.Sub(record.lastTimestamp.Time) if interval > maxInterval { record = aggregateRecord{localKeys: sets.NewString()} } record.localKeys.Insert(localKey) record.lastTimestamp = now // 重新加入到LRU缓存 e.cache.Add(aggregateKey, record) // 如果没有达到阈值则不进行聚合 if uint(record.localKeys.Len()) < e.maxEvents { return newEvent, eventKey } record.localKeys.PopAny() eventCopy := &v1.Event{ ObjectMeta: metav1.ObjectMeta{ Name: fmt.Sprintf("%v.%x", newEvent.InvolvedObject.Name, now.UnixNano()), Namespace: newEvent.Namespace, }, Count: 1, FirstTimestamp: now, InvolvedObject: newEvent.InvolvedObject, LastTimestamp: now, // 否则对message进行聚合 Message: e.messageFunc(newEvent), Type: newEvent.Type, Reason: newEvent.Reason, Source: newEvent.Source, } return eventCopy, aggregateKey }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49EventAggregate首先会去缓存里面查找相同的聚合记录aggregateRecord,如果没有的话,那么会在校验时间间隔的时候顺便创建聚合记录aggregateRecord。由于采用了LRU缓存,所以再将聚合记录重新Add到缓存的头部。接下来判断缓存是否已经超过了阈值,如果没有达到阈值,那么直接返回不进行聚合;否则会重新copy传入的Event,并调用messageFunc方法聚合Message。
# 2.4.4.eventObserve
func (e *eventLogger) eventObserve(newEvent *v1.Event, key string) (*v1.Event, []byte, error) { var ( patch []byte err error ) eventCopy := *newEvent event := &eventCopy e.Lock() defer e.Unlock() // 检查是否在缓存中 lastObservation := e.lastEventObservationFromCache(key) // 如果存在则对count进行自增 if lastObservation.count > 0 { event.Name = lastObservation.name event.ResourceVersion = lastObservation.resourceVersion event.FirstTimestamp = lastObservation.firstTimestamp event.Count = int32(lastObservation.count) + 1 eventCopy2 := *event eventCopy2.Count = 0 eventCopy2.LastTimestamp = metav1.NewTime(time.Unix(0, 0)) eventCopy2.Message = "" newData, _ := json.Marshal(event) oldData, _ := json.Marshal(eventCopy2) patch, err = strategicpatch.CreateTwoWayMergePatch(oldData, newData, event) } // 最后重新更新缓存记录 e.cache.Add( key, eventLog{ count: uint(event.Count), firstTimestamp: event.FirstTimestamp, name: event.Name, resourceVersion: event.ResourceVersion, }, ) return event, patch, err }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42eventObserve会去查找缓存中的记录,然后对count进行自增后更新到缓存。
# 2.4.5.Filter
func (f *EventSourceObjectSpamFilter) Filter(event *v1.Event) bool { var record spamRecord eventKey := f.spamKeyFunc(event) f.Lock() defer f.Unlock() value, found := f.cache.Get(eventKey) if found { record = value.(spamRecord) } if record.rateLimiter == nil { record.rateLimiter = flowcontrol.NewTokenBucketPassiveRateLimiterWithClock(f.qps, f.burst, f.clock) } // 使用令牌桶进行过滤 filter := !record.rateLimiter.TryAccept() // 更新缓存 f.cache.Add(eventKey, record) return filter }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24Filter主要起到一个限速作用,通过令牌桶进行过滤操作。
# 2.4.6.recordEvent
func recordEvent(sink EventSink, event *v1.Event, patch []byte, updateExistingEvent bool, eventCorrelator *EventCorrelator) bool { var newEvent *v1.Event var err error // 更新已经存在的事件 if updateExistingEvent { newEvent, err = sink.Patch(event, patch) } // 创建一个新事件 if !updateExistingEvent || (updateExistingEvent && util.IsKeyNotFoundError(err)) { event.ResourceVersion = "" newEvent, err = sink.Create(event) } if err == nil { eventCorrelator.UpdateState(newEvent) return true } // 如果是已知错误就不要再重试了,否则让上层进行重试 switch err.(type) { case *restclient.RequestConstructionError: return true case *errors.StatusError: return true case *errors.UnexpectedObjectError: // We don't expect this; it implies the server's response didn't match a // known pattern. Go ahead and retry. default: // This case includes actual http transport errors. Go ahead and retry. } return false }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31recordEvent方法会根据eventCorrelator返回的结果决定新建事件还是更新已经存在的事件,并根据请求的解决抉择是否重试。
# 2.5.总结
1.
kubelet首先会初始化EventBroadcaster对象,同时会初始化一个Broadcaster对象2.
kubelet通过EventBroadcaster对象的NewRecorder()方法初始化EventRecorder对象,EventRecorder对象提供的几个方法会生成events并通过ActionOrDrop()方法发送events到Broadcaster的channel队列中3.
Broadcaster的作用就是接收所有的events并进行广播,Broadcaster初始化后会在后台启动一个goroutine,接收所有从EventRecorder发来的events4.
EventBroadcaster对events有三个处理方法:StartEventWatcher()、StartRecordingToSink()、StartStructuredLogging(),StartEventWatcher()是其中的核心方法,会初始化一个watcher注册到Broadcaster,其余两个处理函数对StartEventWatcher()进行了封装,并实现了自己的处理函数。5.
Broadcaster中有一个map会保存每一个注册的watcher,其会将所有的events广播给每一个watcher,每个watcher通过它的ResultChan()方法从channel接收events6.
kubelet会使用StartRecordingToSink()和StartStructuredLogging()对events进行处理,StartRecordingToSink()处理函数收到events后会进行缓存、过滤、聚合而后发送到apiserver,apiserver会将events保存到etcd中,StartStructuredLogging()仅将events保存到kubelet的日志