mirror of
https://github.com/kubernetes/sample-controller.git
synced 2025-02-01 01:12:52 +08:00
Merge pull request #73308 from krzysied/reflector_trace2
Adding trace to reflector initialization Kubernetes-commit: f5f5d9a54a6397012b780e1714abf7d8b4f5037c
This commit is contained in:
parent
da16a67695
commit
abcd26d668
488
Godeps/Godeps.json
generated
488
Godeps/Godeps.json
generated
File diff suppressed because it is too large
Load Diff
14
vendor/k8s.io/client-go/tools/cache/reflector.go
generated
vendored
14
vendor/k8s.io/client-go/tools/cache/reflector.go
generated
vendored
@ -41,6 +41,7 @@ import (
|
|||||||
"k8s.io/apimachinery/pkg/util/wait"
|
"k8s.io/apimachinery/pkg/util/wait"
|
||||||
"k8s.io/apimachinery/pkg/watch"
|
"k8s.io/apimachinery/pkg/watch"
|
||||||
"k8s.io/klog"
|
"k8s.io/klog"
|
||||||
|
"k8s.io/utils/trace"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Reflector watches a specified resource and causes all changes to be reflected in the given store.
|
// Reflector watches a specified resource and causes all changes to be reflected in the given store.
|
||||||
@ -176,6 +177,9 @@ func (r *Reflector) ListAndWatch(stopCh <-chan struct{}) error {
|
|||||||
r.metrics.numberOfLists.Inc()
|
r.metrics.numberOfLists.Inc()
|
||||||
start := r.clock.Now()
|
start := r.clock.Now()
|
||||||
|
|
||||||
|
if err := func() error {
|
||||||
|
initTrace := trace.New("Reflector " + r.name + " ListAndWatch")
|
||||||
|
defer initTrace.LogIfLong(10 * time.Second)
|
||||||
var list runtime.Object
|
var list runtime.Object
|
||||||
var err error
|
var err error
|
||||||
listCh := make(chan struct{}, 1)
|
listCh := make(chan struct{}, 1)
|
||||||
@ -199,22 +203,30 @@ func (r *Reflector) ListAndWatch(stopCh <-chan struct{}) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("%s: Failed to list %v: %v", r.name, r.expectedType, err)
|
return fmt.Errorf("%s: Failed to list %v: %v", r.name, r.expectedType, err)
|
||||||
}
|
}
|
||||||
|
initTrace.Step("Objects listed")
|
||||||
r.metrics.listDuration.Observe(time.Since(start).Seconds())
|
r.metrics.listDuration.Observe(time.Since(start).Seconds())
|
||||||
listMetaInterface, err := meta.ListAccessor(list)
|
listMetaInterface, err := meta.ListAccessor(list)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("%s: Unable to understand list result %#v: %v", r.name, list, err)
|
return fmt.Errorf("%s: Unable to understand list result %#v: %v", r.name, list, err)
|
||||||
}
|
}
|
||||||
resourceVersion = listMetaInterface.GetResourceVersion()
|
resourceVersion = listMetaInterface.GetResourceVersion()
|
||||||
|
initTrace.Step("Resource version extracted")
|
||||||
items, err := meta.ExtractList(list)
|
items, err := meta.ExtractList(list)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("%s: Unable to understand list result %#v (%v)", r.name, list, err)
|
return fmt.Errorf("%s: Unable to understand list result %#v (%v)", r.name, list, err)
|
||||||
}
|
}
|
||||||
|
initTrace.Step("Objects extracted")
|
||||||
r.metrics.numberOfItemsInList.Observe(float64(len(items)))
|
r.metrics.numberOfItemsInList.Observe(float64(len(items)))
|
||||||
if err := r.syncWith(items, resourceVersion); err != nil {
|
if err := r.syncWith(items, resourceVersion); err != nil {
|
||||||
return fmt.Errorf("%s: Unable to sync list result: %v", r.name, err)
|
return fmt.Errorf("%s: Unable to sync list result: %v", r.name, err)
|
||||||
}
|
}
|
||||||
|
initTrace.Step("SyncWith done")
|
||||||
r.setLastSyncResourceVersion(resourceVersion)
|
r.setLastSyncResourceVersion(resourceVersion)
|
||||||
|
initTrace.Step("Resource version updated")
|
||||||
|
return nil
|
||||||
|
}(); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
resyncerrc := make(chan error, 1)
|
resyncerrc := make(chan error, 1)
|
||||||
cancelCh := make(chan struct{})
|
cancelCh := make(chan struct{})
|
||||||
|
96
vendor/k8s.io/utils/trace/trace.go
generated
vendored
Normal file
96
vendor/k8s.io/utils/trace/trace.go
generated
vendored
Normal file
@ -0,0 +1,96 @@
|
|||||||
|
/*
|
||||||
|
Copyright 2015 The Kubernetes Authors.
|
||||||
|
|
||||||
|
Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
you may not use this file except in compliance with the License.
|
||||||
|
You may obtain a copy of the License at
|
||||||
|
|
||||||
|
http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
|
||||||
|
Unless required by applicable law or agreed to in writing, software
|
||||||
|
distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
See the License for the specific language governing permissions and
|
||||||
|
limitations under the License.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package trace
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"fmt"
|
||||||
|
"math/rand"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"k8s.io/klog"
|
||||||
|
)
|
||||||
|
|
||||||
|
type traceStep struct {
|
||||||
|
stepTime time.Time
|
||||||
|
msg string
|
||||||
|
}
|
||||||
|
|
||||||
|
// Trace keeps track of a set of "steps" and allows us to log a specific
|
||||||
|
// step if it took longer than its share of the total allowed time
|
||||||
|
type Trace struct {
|
||||||
|
name string
|
||||||
|
startTime time.Time
|
||||||
|
steps []traceStep
|
||||||
|
}
|
||||||
|
|
||||||
|
// New creates a Trace with the specified name
|
||||||
|
func New(name string) *Trace {
|
||||||
|
return &Trace{name, time.Now(), nil}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Step adds a new step with a specific message
|
||||||
|
func (t *Trace) Step(msg string) {
|
||||||
|
if t.steps == nil {
|
||||||
|
// traces almost always have less than 6 steps, do this to avoid more than a single allocation
|
||||||
|
t.steps = make([]traceStep, 0, 6)
|
||||||
|
}
|
||||||
|
t.steps = append(t.steps, traceStep{time.Now(), msg})
|
||||||
|
}
|
||||||
|
|
||||||
|
// Log is used to dump all the steps in the Trace
|
||||||
|
func (t *Trace) Log() {
|
||||||
|
// an explicit logging request should dump all the steps out at the higher level
|
||||||
|
t.logWithStepThreshold(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *Trace) logWithStepThreshold(stepThreshold time.Duration) {
|
||||||
|
var buffer bytes.Buffer
|
||||||
|
tracenum := rand.Int31()
|
||||||
|
endTime := time.Now()
|
||||||
|
|
||||||
|
totalTime := endTime.Sub(t.startTime)
|
||||||
|
buffer.WriteString(fmt.Sprintf("Trace[%d]: %q (started: %v) (total time: %v):\n", tracenum, t.name, t.startTime, totalTime))
|
||||||
|
lastStepTime := t.startTime
|
||||||
|
for _, step := range t.steps {
|
||||||
|
stepDuration := step.stepTime.Sub(lastStepTime)
|
||||||
|
if stepThreshold == 0 || stepDuration > stepThreshold || klog.V(4) {
|
||||||
|
buffer.WriteString(fmt.Sprintf("Trace[%d]: [%v] [%v] %v\n", tracenum, step.stepTime.Sub(t.startTime), stepDuration, step.msg))
|
||||||
|
}
|
||||||
|
lastStepTime = step.stepTime
|
||||||
|
}
|
||||||
|
stepDuration := endTime.Sub(lastStepTime)
|
||||||
|
if stepThreshold == 0 || stepDuration > stepThreshold || klog.V(4) {
|
||||||
|
buffer.WriteString(fmt.Sprintf("Trace[%d]: [%v] [%v] END\n", tracenum, endTime.Sub(t.startTime), stepDuration))
|
||||||
|
}
|
||||||
|
|
||||||
|
klog.Info(buffer.String())
|
||||||
|
}
|
||||||
|
|
||||||
|
// LogIfLong is used to dump steps that took longer than its share
|
||||||
|
func (t *Trace) LogIfLong(threshold time.Duration) {
|
||||||
|
if time.Since(t.startTime) >= threshold {
|
||||||
|
// if any step took more than it's share of the total allowed time, it deserves a higher log level
|
||||||
|
stepThreshold := threshold / time.Duration(len(t.steps)+1)
|
||||||
|
t.logWithStepThreshold(stepThreshold)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TotalTime can be used to figure out how long it took since the Trace was created
|
||||||
|
func (t *Trace) TotalTime() time.Duration {
|
||||||
|
return time.Since(t.startTime)
|
||||||
|
}
|
Loading…
Reference in New Issue
Block a user