mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
fix(flashnode): record disk dimension to file read and write operations #23131609
Signed-off-by: clinx <chenlin1@oppo.com>
This commit is contained in:
parent
cfa8134cc7
commit
19e4fdc075
@ -26,6 +26,7 @@ import (
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/cubefs/cubefs/flashnode"
|
||||
"github.com/cubefs/cubefs/proto"
|
||||
"github.com/cubefs/cubefs/util"
|
||||
"github.com/cubefs/cubefs/util/auditlog"
|
||||
@ -240,7 +241,6 @@ func (cb *CacheBlock) writeCacheBlockFileHeader(file *os.File) (err error) {
|
||||
return fmt.Errorf("no lru cache item related to dataPath(%v)", cb.rootPath)
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
@ -539,6 +539,10 @@ func (cb *CacheBlock) updateAllocSize(size int64) {
|
||||
cb.allocSize = size
|
||||
}
|
||||
|
||||
func (cb *CacheBlock) GetRootPath() string {
|
||||
return cb.rootPath
|
||||
}
|
||||
|
||||
func (cb *CacheBlock) InitOnceForCacheRead(engine *CacheEngine, sources []*proto.DataSource) {
|
||||
cb.initOnce.Do(func() {
|
||||
cb.InitForCacheRead(sources, engine.readDataNodeTimeout)
|
||||
@ -575,6 +579,8 @@ func (cb *CacheBlock) InitForCacheRead(sources []*proto.DataSource, readDataNode
|
||||
return e
|
||||
}
|
||||
offset += size
|
||||
flashnode.UpdateWriteBytesMetric(uint64(size), cb.GetRootPath())
|
||||
flashnode.UpdateWriteCountMetric(cb.GetRootPath())
|
||||
return nil
|
||||
}
|
||||
logPrefix := func() string {
|
||||
|
||||
@ -163,7 +163,7 @@ func doStart(s common.Server, cfg *config.Config) (err error) {
|
||||
if err = f.start(cfg); err != nil {
|
||||
return
|
||||
}
|
||||
f.registerMetrics()
|
||||
f.registerMetrics(f.disks)
|
||||
exporter.RegistConsul(f.clusterID, moduleName, cfg)
|
||||
f.startMetrics()
|
||||
return
|
||||
|
||||
@ -224,8 +224,8 @@ func (f *FlashNode) doStreamReadRequest(ctx context.Context, conn net.Conn, req
|
||||
if err != nil {
|
||||
log.LogWarnf("%s cache block(%v) err:%v", action, block.String(), err)
|
||||
} else {
|
||||
f.updateReadCountMetric()
|
||||
f.updateReadBytesMetric(req.Size_)
|
||||
UpdateReadCountMetric(block.GetRootPath())
|
||||
UpdateReadBytesMetric(req.Size_, block.GetRootPath())
|
||||
}
|
||||
}()
|
||||
for needReplySize > 0 {
|
||||
|
||||
@ -1,6 +1,8 @@
|
||||
package flashnode
|
||||
|
||||
import (
|
||||
"github.com/cubefs/cubefs/flashnode/cachengine"
|
||||
"path"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
@ -12,13 +14,19 @@ const (
|
||||
StatPeriod = time.Minute * time.Duration(1)
|
||||
MetricFlashNodeReadBytes = "flashNodeReadBytes"
|
||||
MetricFlashNodeReadCount = "flashNodeReadCount"
|
||||
MetricFlashNodeWriteBytes = "flashNodeWriteBytes"
|
||||
MetricFlashNodeWriteCount = "flashNodeWriteCount"
|
||||
MetricFlashNodeHitRate = "flashNodeHitRate"
|
||||
MetricFlashNodeEvictCount = "flashNodeEvictCount"
|
||||
)
|
||||
|
||||
var StatMap = make(map[string]*MetricStat)
|
||||
|
||||
type MetricStat struct {
|
||||
ReadBytes uint64
|
||||
ReadCount uint64
|
||||
ReadBytes uint64
|
||||
ReadCount uint64
|
||||
WriteBytes uint64
|
||||
WriteCount uint64
|
||||
}
|
||||
|
||||
type FlashNodeMetrics struct {
|
||||
@ -26,20 +34,28 @@ type FlashNodeMetrics struct {
|
||||
stopC chan struct{}
|
||||
MetricReadBytes *exporter.Gauge
|
||||
MetricReadCount *exporter.Gauge
|
||||
MetricWriteBytes *exporter.Gauge
|
||||
MetricWriteCount *exporter.Gauge
|
||||
MetricEvictCount *exporter.Gauge
|
||||
MetricHitRate *exporter.Gauge
|
||||
Stat MetricStat
|
||||
}
|
||||
|
||||
func (f *FlashNode) registerMetrics() {
|
||||
func (f *FlashNode) registerMetrics(disks []*cachengine.Disk) {
|
||||
f.metrics = &FlashNodeMetrics{
|
||||
flashNode: f,
|
||||
stopC: make(chan struct{}),
|
||||
}
|
||||
|
||||
f.metrics.MetricReadBytes = exporter.NewGauge(MetricFlashNodeReadBytes)
|
||||
f.metrics.MetricReadCount = exporter.NewGauge(MetricFlashNodeReadCount)
|
||||
f.metrics.MetricWriteBytes = exporter.NewGauge(MetricFlashNodeWriteBytes)
|
||||
f.metrics.MetricWriteCount = exporter.NewGauge(MetricFlashNodeWriteCount)
|
||||
f.metrics.MetricEvictCount = exporter.NewGauge(MetricFlashNodeEvictCount)
|
||||
f.metrics.MetricHitRate = exporter.NewGauge(MetricFlashNodeHitRate)
|
||||
for _, d := range disks {
|
||||
StatMap[path.Join(d.Path, cachengine.DefaultCacheDirName)] = new(MetricStat)
|
||||
}
|
||||
|
||||
log.LogInfof("registerMetrics")
|
||||
}
|
||||
|
||||
@ -67,18 +83,38 @@ func (fm *FlashNodeMetrics) doStat() {
|
||||
log.LogInfof("FlashNodeMetrics: doStat")
|
||||
fm.setReadBytesMetric()
|
||||
fm.setReadCountMetric()
|
||||
fm.setWriteBytesMetric()
|
||||
fm.setWriteCountMetric()
|
||||
fm.setEvictCountMetric()
|
||||
fm.setHitRateMetric()
|
||||
}
|
||||
|
||||
func (fm *FlashNodeMetrics) setReadBytesMetric() {
|
||||
readBytes := atomic.SwapUint64(&fm.Stat.ReadBytes, 0)
|
||||
fm.MetricReadBytes.SetWithLabels(float64(readBytes), map[string]string{"cluster": fm.flashNode.clusterID, exporter.FlashNode: fm.flashNode.localAddr})
|
||||
for d, stat := range StatMap {
|
||||
readBytes := atomic.SwapUint64(&stat.ReadBytes, 0)
|
||||
fm.MetricReadBytes.SetWithLabels(float64(readBytes), map[string]string{"cluster": fm.flashNode.clusterID, exporter.FlashNode: fm.flashNode.localAddr, exporter.Disk: d})
|
||||
}
|
||||
}
|
||||
|
||||
func (fm *FlashNodeMetrics) setReadCountMetric() {
|
||||
readCount := atomic.SwapUint64(&fm.Stat.ReadCount, 0)
|
||||
fm.MetricReadCount.SetWithLabels(float64(readCount), map[string]string{"cluster": fm.flashNode.clusterID, exporter.FlashNode: fm.flashNode.localAddr})
|
||||
for d, stat := range StatMap {
|
||||
readCount := atomic.SwapUint64(&stat.ReadCount, 0)
|
||||
fm.MetricReadCount.SetWithLabels(float64(readCount), map[string]string{"cluster": fm.flashNode.clusterID, exporter.FlashNode: fm.flashNode.localAddr, exporter.Disk: d})
|
||||
}
|
||||
}
|
||||
|
||||
func (fm *FlashNodeMetrics) setWriteBytesMetric() {
|
||||
for d, stat := range StatMap {
|
||||
writeBytes := atomic.SwapUint64(&stat.WriteBytes, 0)
|
||||
fm.MetricWriteBytes.SetWithLabels(float64(writeBytes), map[string]string{"cluster": fm.flashNode.clusterID, exporter.FlashNode: fm.flashNode.localAddr, exporter.Disk: d})
|
||||
}
|
||||
}
|
||||
|
||||
func (fm *FlashNodeMetrics) setWriteCountMetric() {
|
||||
for d, stat := range StatMap {
|
||||
writeCount := atomic.SwapUint64(&stat.WriteCount, 0)
|
||||
fm.MetricWriteCount.SetWithLabels(float64(writeCount), map[string]string{"cluster": fm.flashNode.clusterID, exporter.FlashNode: fm.flashNode.localAddr, exporter.Disk: d})
|
||||
}
|
||||
}
|
||||
|
||||
func (fm *FlashNodeMetrics) setEvictCountMetric() {
|
||||
@ -95,14 +131,26 @@ func (fm *FlashNodeMetrics) setHitRateMetric() {
|
||||
}
|
||||
}
|
||||
|
||||
func (f *FlashNode) updateReadBytesMetric(size uint64) {
|
||||
if f.metrics != nil {
|
||||
atomic.AddUint64(&f.metrics.Stat.ReadBytes, size)
|
||||
func UpdateReadBytesMetric(size uint64, d string) {
|
||||
if stat, ok := StatMap[d]; ok {
|
||||
atomic.AddUint64(&stat.ReadBytes, size)
|
||||
}
|
||||
}
|
||||
|
||||
func (f *FlashNode) updateReadCountMetric() {
|
||||
if f.metrics != nil {
|
||||
atomic.AddUint64(&f.metrics.Stat.ReadCount, 1)
|
||||
func UpdateReadCountMetric(d string) {
|
||||
if stat, ok := StatMap[d]; ok {
|
||||
atomic.AddUint64(&stat.ReadCount, 1)
|
||||
}
|
||||
}
|
||||
|
||||
func UpdateWriteBytesMetric(size uint64, d string) {
|
||||
if stat, ok := StatMap[d]; ok {
|
||||
atomic.AddUint64(&stat.WriteBytes, size)
|
||||
}
|
||||
}
|
||||
|
||||
func UpdateWriteCountMetric(d string) {
|
||||
if stat, ok := StatMap[d]; ok {
|
||||
atomic.AddUint64(&stat.WriteCount, 1)
|
||||
}
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user