()
| 49 | func (sd *StatsDumper) SetWriter(w io.Writer) { sd.w = w } |
| 50 | |
| 51 | func (sd *StatsDumper) dumpNow() { |
| 52 | sd.lock.Lock() |
| 53 | defer sd.lock.Unlock() |
| 54 | |
| 55 | t := sd.t |
| 56 | nsec := int64(time.Now().UTC().Sub(sd.start).Seconds()) |
| 57 | |
| 58 | istats := t.Input.Stats() |
| 59 | currlines := istats.NumProcessedLines |
| 60 | |
| 61 | // Collect metrics from input, filters and outputs that we can |
| 62 | // forward to statsd |
| 63 | allMetrics := make(MetricsBag) |
| 64 | allMetrics.Merge(istats.Metrics) |
| 65 | |
| 66 | var filtered int64 |
| 67 | filteredMap := make(map[string]int64) |
| 68 | for fidx, f := range t.Filters { |
| 69 | stats := f.Stats() |
| 70 | if stats.NumFilteredLines > 0 { |
| 71 | sd.metrics.RawCountWithTags("filtered_lines", stats.NumFilteredLines, sd.filterTags[fidx]) |
| 72 | filtered += stats.NumFilteredLines |
| 73 | filteredMap[fmt.Sprintf("%T", f)] += filtered |
| 74 | } |
| 75 | allMetrics.Merge(stats.Metrics) |
| 76 | } |
| 77 | |
| 78 | var curwlines, outErrors int64 |
| 79 | for _, o := range t.Output { |
| 80 | stats := o.Stats() |
| 81 | outErrors += stats.NumErrorLines |
| 82 | curwlines += stats.NumProcessedLines |
| 83 | allMetrics.Merge(stats.Metrics) |
| 84 | } |
| 85 | sd.metrics.RawCount("processed_lines", curwlines) |
| 86 | |
| 87 | var numUploadErrors, numUploads int64 |
| 88 | if t.Upload != nil { |
| 89 | uStats := t.Upload.Stats() |
| 90 | numUploads = uStats.NumProcessedFiles |
| 91 | numUploadErrors = uStats.NumErrorFiles |
| 92 | sd.metrics.RawCount("uploads", numUploads) |
| 93 | sd.metrics.RawCount("upload_errors", numUploadErrors) |
| 94 | allMetrics.Merge(uStats.Metrics) |
| 95 | } |
| 96 | |
| 97 | if numUploads < sd.prevUploads { |
| 98 | log.Fatalf("numUploads < prevUploads: %d < %d\n", numUploads, sd.prevUploads) |
| 99 | } |
| 100 | |
| 101 | invalid := sd.countInvalid() |
| 102 | parseErrors := t.malformed |
| 103 | totalErrors := invalid + parseErrors + filtered + outErrors |
| 104 | sd.metrics.RawCount("error_lines", totalErrors) |
| 105 | |
| 106 | for k, v := range allMetrics { |
| 107 | switch k[0] { |
| 108 | case 'c': |
no test coverage detected