数据库其它

关注公众号 jb51net

关闭
首页 > 数据库 > 数据库其它 > VictoriaMetrics写入索引

时序数据库VictoriaMetrics源码解析之写入与索引

作者:a朋

这篇文章主要为大家介绍了VictoriaMetrics时序数据库的写入与索引源码解析,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪

一. 存储格式

下图是向VictoriaMetrics写入prometheus协议数据的示例:

VM在收到写入请求时,会对请求中包含的时序数据做转换处理:

因此,VM的数据整体上分为索引和数据2个部分:

二. 整体流程

VictoriaMetrics在写入原始的rows数据时,写入过程分为两个部分:

写入流程:

三. 写入代码

1.入口代码

vmstorage监听tcp端口,收到vminsert的插入请求后,进行处理:

// app/vmstorage/servers/vminsert.go
func (s *VMInsertServer) run() {
    ...
    for {
        c, err := s.ln.Accept()
        ...
        go func() {
            bc, err := handshake.VMInsertServer(c, compressionLevel)
            ...
            err = clusternative.ParseStream(bc, func(rows []storage.MetricRow) error {
                vminsertMetricsRead.Add(len(rows))
                return s.storage.AddRows(rows, uint8(*precisionBits))    // 入口代码
            }, s.storage.IsReadOnly)
            ...
        }()
    }
}

写入时,1次最多写8K个rows:

func (s *Storage) AddRows(mrs []MetricRow, precisionBits uint8) error {
    ....
    maxBlockLen := len(ic.rrs)
    for len(mrs) > 0 {
        mrsBlock := mrs
        // 一次最多写8K,maxBlockLen=8000
        if len(mrs) > maxBlockLen {
            mrsBlock = mrs[:maxBlockLen]
            mrs = mrs[maxBlockLen:]
        } else {
            mrs = nil
        }
        // 写入8K rows的数据
        if err := s.add(ic.rrs, ic.tmpMrs, mrsBlock, precisionBits); err != nil {
            if firstErr == nil {
                firstErr = err
            }
            continue
        }
        atomic.AddUint64(&rowsAddedTotal, uint64(len(mrsBlock)))
    }
    ....
}

2.写入流程的代码

写入过程主要分2步:

// lib/storage/storage.go
func (s *Storage) add(rows []rawRow, dstMrs []*MetricRow, mrs []MetricRow, precisionBits uint8) error {
    ...
    // 1.构造r.TSID
    // 若跟prevMetricNameRaw相同,则使用pervTSID;
    // 若cache中有metricNameRaw,则使用cache.TSID;
    for i := range mrs {
        mr := &mrs[i]
        ...
        dstMrs[j] = mr
        r := &rows[j]
        j++
        r.Timestamp = mr.Timestamp
        r.Value = mr.Value
        r.PrecisionBits = precisionBits
        if string(mr.MetricNameRaw) == string(prevMetricNameRaw) {    // 使用prevTSID
            // Fast path - the current mr contains the same metric name as the previous mr, so it contains the same TSID.
            // This path should trigger on bulk imports when many rows contain the same MetricNameRaw.
            r.TSID = prevTSID
            continue
        }
        if s.getTSIDFromCache(&genTSID, mr.MetricNameRaw) {        // 使用缓存的TSID
            ...
            r.TSID = genTSID.TSID
            prevTSID = r.TSID
            prevMetricNameRaw = mr.MetricNameRaw
            ...
            continue
        }
        ...
    }
    if pmrs != nil {
        // Sort pendingMetricRows by canonical metric name in order to speed up search via `is` in the loop below.
        pendingMetricRows := pmrs.pmrs
        sort.Slice(pendingMetricRows, func(i, j int) bool {
            return string(pendingMetricRows[i].MetricName) < string(pendingMetricRows[j].MetricName)
        })
        prevMetricNameRaw = nil
        var slowInsertsCount uint64
        for i := range pendingMetricRows {
            ...
            r := &rows[j]
            j++
            r.Timestamp = mr.Timestamp
            r.Value = mr.Value
            r.PrecisionBits = precisionBits
            // 尝试去index找查找,或者创建
          if err := is.GetOrCreateTSIDByName(&r.TSID, pmr.MetricName, mr.MetricNameRaw, date); err != nil {
                ...
                continue
            }
            genTSID.generation = idb.generation
            genTSID.TSID = r.TSID
            // 放回cache
            s.putTSIDToCache(&genTSID, mr.MetricNameRaw)
            prevTSID = r.TSID
            prevMetricNameRaw = mr.MetricNameRaw
        }
    }
    ...
    dstMrs = dstMrs[:j]
    rows = rows[:j]
    err := s.updatePerDateData(rows, dstMrs)
    if err != nil {
        err = fmt.Errorf("cannot update per-date data: %w", err)
    } else {
        // TSID构造完毕,开始插入数据
        err = s.tb.AddRows(rows)
        ...
    }
    ...
    return nil
}

3.写index

写index是slow path,重点看一下:

// lib/storage/index_db.go
func (is *indexSearch) GetOrCreateTSIDByName(dst *TSID, metricName, metricNameRaw []byte, date uint64) error {
    // 1.首先尝试在index中查找
    if is.tsidByNameMisses < 100 {
        err := is.getTSIDByMetricName(dst, metricName)
        // 在index中找到了
        if err == nil {
            // Fast path - the TSID for the given metricName has been found in the index.
            is.tsidByNameMisses = 0
            if err = is.db.s.registerSeriesCardinality(dst.MetricID, metricNameRaw); err != nil {
                return err
            }
            return nil
        }
        is.tsidByNameMisses++
    } else {
        is.tsidByNameSkips++
        if is.tsidByNameSkips > 10000 {
            is.tsidByNameSkips = 0
            is.tsidByNameMisses = 0
        }
    }
    // 2.没有找到,那么创建一个
    if err := is.createTSIDByName(dst, metricName, metricNameRaw, date); err != nil {
        userReadableMetricName := getUserReadableMetricName(metricNameRaw)
        return fmt.Errorf("cannot create TSID by MetricName %s: %w", userReadableMetricName, err)
    }
    return nil
}

4. 生成TSID

具体生成TSID的逻辑:

// lib/storage/index_db.go
func generateTSID(dst *TSID, mn *MetricName) {
    dst.AccountID = mn.AccountID
    dst.ProjectID = mn.ProjectID
    dst.MetricGroupID = xxhash.Sum64(mn.MetricGroup)
    if len(mn.Tags) > 0 {
        dst.JobID = uint32(xxhash.Sum64(mn.Tags[0].Value))
    }
    if len(mn.Tags) > 1 {
        dst.InstanceID = uint32(xxhash.Sum64(mn.Tags[1].Value))
    }
    dst.MetricID = generateUniqueMetricID()
}

而TSID中的metricID是由启动时的时间戳+1产生:

// Returns local unique MetricID.
func generateUniqueMetricID() uint64 {
    return atomic.AddUint64(&amp;nextUniqueMetricID, 1)
}
var nextUniqueMetricID = uint64(time.Now().UnixNano())

5. 创建index items

// lib/storage/index_db.go
func (is *indexSearch) createGlobalIndexes(tsid *TSID, mn *MetricName) {
    // The order of index items is important.
    // It guarantees index consistency.
    ii := getIndexItems()
    defer putIndexItems(ii)
    // Create MetricName -> TSID index.
    ii.B = append(ii.B, nsPrefixMetricNameToTSID)
    ii.B = mn.Marshal(ii.B)
    ii.B = append(ii.B, kvSeparatorChar)
    ii.B = tsid.Marshal(ii.B)
    ii.Next()
    // Create MetricID -> MetricName index.
    ii.B = marshalCommonPrefix(ii.B, nsPrefixMetricIDToMetricName, mn.AccountID, mn.ProjectID)
    ii.B = encoding.MarshalUint64(ii.B, tsid.MetricID)
    ii.B = mn.Marshal(ii.B)
    ii.Next()
    // Create MetricID -> TSID index.
    ii.B = marshalCommonPrefix(ii.B, nsPrefixMetricIDToTSID, mn.AccountID, mn.ProjectID)
    ii.B = encoding.MarshalUint64(ii.B, tsid.MetricID)
    ii.B = tsid.Marshal(ii.B)
    ii.Next()
    prefix := kbPool.Get()
    prefix.B = marshalCommonPrefix(prefix.B[:0], nsPrefixTagToMetricIDs, mn.AccountID, mn.ProjectID)
    ii.registerTagIndexes(prefix.B, mn, tsid.MetricID)
    kbPool.Put(prefix)
    is.db.tb.AddItems(ii.Items)     // 将items存入内存shards
}

6. index items存入内存shards

Index items构造完成后,被写入内存的shards,会有异步的goroutine将其压缩写入disk。

写内存shards的方法: roundRobin

// lib/mergeset/table.go
func (riss *rawItemsShards) addItems(tb *Table, items [][]byte) {
   shards := riss.shards
   shardsLen := uint32(len(shards))
   for len(items) > 0 {
      n := atomic.AddUint32(&riss.shardIdx, 1)
      idx := n % shardsLen
      items = shards[idx].addItems(tb, items)
   }
}

内存中shards总数,跟cpu核数有关系:

// lib/mergeset/table.go
/ The number of shards for rawItems per table.
//
// Higher number of shards reduces CPU contention and increases the max bandwidth on multi-core systems.
var rawItemsShardsPerTable = func() int {
   cpus := cgroup.AvailableCPUs()
   multiplier := cpus
   if multiplier > 16 {
      multiplier = 16
   }
   return (cpus*multiplier + 1) / 2
}()

以上就是时序数据库VictoriaMetrics源码解析之写入与索引的详细内容,更多关于VictoriaMetrics写入索引的资料请关注脚本之家其它相关文章!

您可能感兴趣的文章:
阅读全文