-
Notifications
You must be signed in to change notification settings - Fork 0
/
event_loki.go
115 lines (103 loc) · 2.88 KB
/
event_loki.go
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
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
/**
* Tencent is pleased to support the open source community by making Polaris available.
*
* Copyright (C) 2019 THL A29 Limited, a Tencent company. All rights reserved.
*
* Licensed under the BSD 3-Clause License (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://opensource.org/licenses/BSD-3-Clause
*
* 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 loki
import (
"time"
"github.com/polarismesh/polaris/common/model"
"github.com/polarismesh/polaris/plugin"
)
const (
PluginName = "discoverEventLoki"
defaultBatchSize = 512
defaultQueueSize = 1024
)
type discoverEventLoki struct {
eventCh chan model.DiscoverEvent
stopCh chan struct{}
eventLog *LokiLogger
}
// Name 插件名称
// @return string 返回插件名称
func (d *discoverEventLoki) Name() string {
return PluginName
}
// Initialize 根据配置文件进行初始化插件 discoverEventLoki
// @param conf 配置文件内容
// @return error 初始化失败,返回 error 信息
func (d *discoverEventLoki) Initialize(conf *plugin.ConfigEntry) error {
var queueSize = defaultQueueSize
if val, ok := conf.Option["queueSize"]; ok {
queueSize, _ = val.(int)
}
lokiLogger, err := newLokiLogger(conf.Option)
if err != nil {
return err
}
d.eventLog = lokiLogger
d.eventCh = make(chan model.DiscoverEvent, queueSize)
d.stopCh = make(chan struct{}, 1)
go d.Run()
return nil
}
// Destroy 执行插件销毁
func (d *discoverEventLoki) Destroy() error {
close(d.stopCh)
return nil
}
// PublishEvent 发布一个服务事件
func (d *discoverEventLoki) PublishEvent(event model.DiscoverEvent) {
select {
case d.eventCh <- event:
return
default:
// do nothing
}
}
// Run 执行主逻辑
func (d *discoverEventLoki) Run() {
// 定时刷新事件到日志的定时器
syncInterval := time.NewTicker(time.Duration(10) * time.Second)
defer syncInterval.Stop()
batch := make([]model.DiscoverEvent, 0, defaultBatchSize)
batchSize := 0
for {
select {
case event := <-d.eventCh:
// 确保事件是顺序的
event.CreateTime = time.Now()
batch = append(batch, event)
batchSize++
// 触发批量生产发送 log 阈值
if batchSize == defaultBatchSize {
d.eventLog.Log(batch[:batchSize])
batch = make([]model.DiscoverEvent, 0, defaultBatchSize)
batchSize = 0
}
case <-syncInterval.C:
if batchSize > 0 {
d.eventLog.Log(batch[:batchSize])
batch = make([]model.DiscoverEvent, 0, defaultBatchSize)
batchSize = 0
}
case <-d.stopCh:
if batchSize > 0 {
d.eventLog.Log(batch[:batchSize])
}
return
}
}
}