远程写入prometheus存储具体方法

prometheus一般都是采用pull方式获取数据,但是有一些情况下,不方便配置exporter,就希望能通过push的方式上传指标数据。

简介

prometheus一般都是采用pull方式获取数据,但是有一些情况下,不方便配置exporter,就希望能通过push的方式上传指标数据。

1、可以采用pushgateway的方式,推送到pushgateway,然后prometheus通过pushgateway拉取数据。

2、在新版本中增加了一个参数:–enable-feature=remote-write-receiver,允许远程通过接口/api/v1/write,直接写数据到prometheus里面。

pushgateway在高并发的情况下还是比较消耗资源的,特别是开启一致性检查,高并发写入的时候特别慢。

第二种方式少了一层转发,速度应该比较快。

接口

可以通过prometheus的http接口/api/v1/write提交数据,这个接口的数据格式有有要求: 使用POST方式提交 需要经过protobuf编码,依赖github.com/gogo/protobuf/proto 可以使用snappy进行压缩,依赖github.com/golang/snappy

步骤:

收集指标名称,时间戳,值和标签 将数据转换成prometheus需要的数据格式 使用proto对数据进行编码,并用snappy进行压缩 通过httpClient提交数据

package prome

import (
   "bufio"
   "bytes"
   "context"
   "io"
   "io/ioutil"
   "net/http"
   "net/url"
   "regexp"
   "time"

   "github.com/gogo/protobuf/proto"
   "github.com/golang/snappy"
   "github.com/opentracing-contrib/go-stdlib/nethttp"
   opentracing "github.com/opentracing/opentracing-go"
   "github.com/pkg/errors"
   "github.com/prometheus/common/model"
   "github.com/prometheus/prometheus/pkg/labels"
   "github.com/prometheus/prometheus/prompb"
)

type RecoverableError struct {
   error
}

type HttpClient struct {
   url     *url.URL
   Client  *http.Client
   timeout time.Duration
}

var MetricNameRE = regexp.MustCompile(`^[a-zA-Z_:][a-zA-Z0-9_:]*$`)

type MetricPoint struct {
   Metric  string            `json:"metric"` // 指标名称
   TagsMap map[string]string `json:"tags"`   // 数据标签
   Time    int64             `json:"time"`   // 时间戳,单位是秒
   Value   float64           `json:"value"`  // 内部字段,最终转换之后的float64数值
}

func (c *HttpClient) remoteWritePost(req []byte) error {
   httpReq, err := http.NewRequest("POST", c.url.String(), bytes.NewReader(req))
   if err != nil {
       return err
   }
   httpReq.Header.Add("Content-Encoding""snappy")
   httpReq.Header.Set("Content-Type""application/x-protobuf")
   httpReq.Header.Set("User-Agent""opcai")
   httpReq.Header.Set("X-Prometheus-Remote-Write-Version""0.1.0")
   ctx, cancel := context.WithTimeout(context.Background(), c.timeout)
   defer cancel()

   httpReq = httpReq.WithContext(ctx)

   if parentSpan := opentracing.SpanFromContext(ctx); parentSpan != nil {
       var ht *nethttp.Tracer
       httpReq, ht = nethttp.TraceRequest(
           parentSpan.Tracer(),
           httpReq,
           nethttp.OperationName("Remote Store"),
           nethttp.ClientTrace(false),
       )
       defer ht.Finish()
   }

   httpResp, err := c.Client.Do(httpReq)
   if err != nil {
       // Errors from Client.Do are from (for example) network errors, so are
       // recoverable.
       return RecoverableError{err}
   }
   defer func() {
       io.Copy(ioutil.Discard, httpResp.Body)
       httpResp.Body.Close()
   }()

   if httpResp.StatusCode/100 != 2 {
       scanner := bufio.NewScanner(io.LimitReader(httpResp.Body, 512))
       line := ""
       if scanner.Scan() {
           line = scanner.Text()
       }
       err = errors.Errorf("server returned HTTP status %s: %s", httpResp.Status, line)
   }
   if httpResp.StatusCode/100 == 5 {
       return RecoverableError{err}
   }
   return err
}

func buildWriteRequest(samples []*prompb.TimeSeries) ([]byte, error) {

   req := &prompb.WriteRequest{
       Timeseries: samples,
   }
   data, err := proto.Marshal(req)
   if err != nil {
       return nil, err
   }
   compressed := snappy.Encode(nil, data)
   return compressed, nil
}

type sample struct {
   labels labels.Labels
   t      int64
   v      float64
}

const (
   LABEL_NAME = "__name__"
)

func convertOne(item *MetricPoint) (*prompb.TimeSeries, error) {
   pt := prompb.TimeSeries{}
   pt.Samples = []prompb.Sample{{}}
   s := sample{}
   s.t = item.Time
   s.v = item.Value
   // name
   if !MetricNameRE.MatchString(item.Metric) {
       return &pt, errors.New("invalid metrics name")
   }
   nameLs := labels.Label{
       Name:  LABEL_NAME,
       Value: item.Metric,
   }
   s.labels = append(s.labels, nameLs)
   for k, v := range item.TagsMap {
       if model.LabelNameRE.MatchString(k) {
           ls := labels.Label{
               Name:  k,
               Value: v,
           }
           s.labels = append(s.labels, ls)
       }
   }

   pt.Labels = labelsToLabelsProto(s.labels, pt.Labels)
   // 时间赋值问题,使用毫秒时间戳
   tsMs := time.Unix(s.t, 0).UnixNano() / 1e6
   pt.Samples[0].Timestamp = tsMs
   pt.Samples[0].Value = s.v
   return &pt, nil
}

func labelsToLabelsProto(labels labels.Labels, buf []*prompb.Label) []*prompb.Label {
   result := buf[:0]
   if cap(buf) for _, l := range labels {
       result = append(result, &prompb.Label{
           Name:  l.Name,
           Value: l.Value,
       })
   }
   return result
}

func (c *HttpClient) RemoteWrite(items []MetricPoint) (err error) {
   if len(items) == 0 {
       return
   }
   ts := make([]*prompb.TimeSeries, len(items))
   for i := range items {
       ts[i], err = convertOne(&items[i])
       if err != nil {
           return
       }
   }
   data, err := buildWriteRequest(ts)
   if err != nil {
       return
   }
   err = c.remoteWritePost(data)
   return
}

func NewClient(ur string, timeout time.Duration) (c *HttpClient, err error) {
   u, err := url.Parse(ur)
   if err != nil {
       return
   }
   c = &HttpClient{
       url:     u,
       Client:  &http.Client{},
       timeout: timeout,
   }
   return
}

测试

prometheus启动的时候记得加参数–enable-feature=remote-write-receiver

package prome

import (
   "testing"
   "time"
)

func TestRemoteWrite(t *testing.T) {
   c, err := NewClient("http://localhost:9090/api/v1/write", 10*time.Second)
   if err != nil {
       t.Fatal(err)
   }
   metrics := []MetricPoint{
       {Metric: "opcai1",
           TagsMap: map[string]string{"env""testing""op""opcai"},
           Time:    time.Now().Add(-1 * time.Minute).Unix(),
           Value:   1},
       {Metric: "opcai2",
           TagsMap: map[string]string{"env""testing""op""opcai"},
           Time:    time.Now().Add(-2 * time.Minute).Unix(),
           Value:   2},
       {Metric: "opcai3",
           TagsMap: map[string]string{"env""testing""op""opcai"},
           Time:    time.Now().Unix(),
           Value:   3},
       {Metric: "opcai4",
           TagsMap: map[string]string{"env""testing""op""opcai"},
           Time:    time.Now().Unix(),
           Value:   4},
   }
   err = c.RemoteWrite(metrics)
   if err != nil {
       t.Fatal(err)
   }
   t.Log("end...")
}

使用go test进行测试

go test -v

文章来源网络,作者:管理,如若转载,请注明出处:https://shuyeidc.com/wp/220585.html<

(0)
管理的头像管理
上一篇2025-04-14 15:36
下一篇 2025-04-14 15:37

相关推荐

  • jsp空间购买和交换数据空间怎么买,有哪些注意事项?

    购买JSP空间时,是否考虑过数据交换空间的性能?简米科技(2003年始创,23年行业沉淀)与酷番云(工信部一类增值电信全牌照)这类持牌自营机房的服务商,能确保数据交换的高效稳定,是值得优先选择的合作伙伴,为什么JSP空间需要搭配独立的数据交换空间从JSP应用特性看数据交换需求JSP基于Java技术,常用于企业级……

    2026-08-11
    0
  • 建网站用香港空间效果怎么样,香港空间稳定吗?

    建网站用香港空间,对于创建网站资产来说,核心价值在于免备案和全球带宽优势,尤其适合外贸、跨境电商和需要快速启动的项目,但你必须权衡国内访问延迟,并选择有资质的服务商以保证资产安全,香港空间的核心优势与适用边界免备案:节省时间就是节省成本国内服务器需要备案,通常需要10到20天,香港空间无需备案,域名解析后即可上……

    2026-08-11
    0
  • Java连接云数据库的方法是什么,如何操作

    Java连接云数据库的核心在于通过JDBC驱动,结合云服务商提供的连接地址、端口、数据库名及认证信息,配置安全策略(如SSL、IP白名单),即可实现稳定高效的远程数据库访问,基础准备:JDBC驱动与依赖管理连接云数据库前,需要确保开发环境具备对应的JDBC驱动,以最常见的MySQL为例,你需要引入mysql-c……

    2026-08-11
    0
  • 建网站公安联网备案必须使用数据码吗,备案流程是什么

    网站备案包括ICP备案和公安联网备案,两者缺一不可,公安联网备案必须使用服务商提供的数据码,选择持有合法资质的服务商是顺利通过备案的前提,为什么网站必须进行公安联网备案根据公安部《计算机信息网络国际联网安全保护管理办法》,网站开通后30日内必须到公安机关办理备案手续,未完成公安备案的网站,面临责令整改、关闭网站……

    2026-08-10
    0
  • 建一个企业网站大概需要多少钱?,怎么收费?

    建网站要多少钱,没有一个固定的数字,几百到几万都可能,但真正的“创建网站资产”绝不仅仅是初次投入的成本,而是基于长期稳定、合规和安全的持续性投入,其中核心取决于你选择了什么样的“地基”来承载你的业务,建站预算的构成与行业基准当你开始规划一个网站,最先面对的就是预算问题,一个常见的误区是只关注网站“看起来”的建造……

    2026-08-10
    0

发表回复

您的邮箱地址不会被公开。必填项已用 * 标注