2020-02-23 12:35:47 +01:00
package remotewrite
import (
"flag"
"fmt"
"sync/atomic"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/flagutil"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/httpserver"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/logger"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/memory"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/persistentqueue"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/prompbmarshal"
2020-03-03 12:08:17 +01:00
"github.com/VictoriaMetrics/VictoriaMetrics/lib/promrelabel"
2020-02-23 12:35:47 +01:00
"github.com/VictoriaMetrics/metrics"
xxhash "github.com/cespare/xxhash/v2"
)
var (
remoteWriteURLs = flagutil . NewArray ( "remoteWrite.url" , "Remote storage URL to write data to. It must support Prometheus remote_write API. " +
"It is recommended using VictoriaMetrics as remote storage. Example url: http://<victoriametrics-host>:8428/api/v1/write . " +
"Pass multiple -remoteWrite.url flags in order to write data concurrently to multiple remote storage systems" )
2020-03-03 12:08:17 +01:00
relabelConfigPaths = flagutil . NewArray ( "remoteWrite.urlRelabelConfig" , "Optional path to relabel config for the corresponding -remoteWrite.url" )
tmpDataPath = flag . String ( "remoteWrite.tmpDataPath" , "vmagent-remotewrite-data" , "Path to directory where temporary data for remote write component is stored" )
queues = flag . Int ( "remoteWrite.queues" , 1 , "The number of concurrent queues to each -remoteWrite.url. Set more queues if a single queue " +
2020-02-23 12:35:47 +01:00
"isn't enough for sending high volume of collected data to remote storage" )
showRemoteWriteURL = flag . Bool ( "remoteWrite.showURL" , false , "Whether to show -remoteWrite.url in the exported metrics. " +
"It is hidden by default, since it can contain sensistive auth info" )
2020-03-03 18:48:46 +01:00
maxPendingBytesPerURL = flag . Int ( "remoteWrite.maxDiskUsagePerURL" , 0 , "The maximum file-based buffer size in bytes at -remoteWrite.tmpDataPath " +
"for each -remoteWrite.url. When buffer size reaches the configured maximum, then old data is dropped when adding new data to the buffer. " +
"Buffered data is stored in ~500MB chunks, so the minimum practical value for this flag is 500000000. " +
"Disk usage is unlimited if the value is set to 0" )
2020-02-23 12:35:47 +01:00
)
2020-03-03 12:08:17 +01:00
var rwctxs [ ] * remoteWriteCtx
2020-02-23 12:35:47 +01:00
// Init initializes remotewrite.
//
// It must be called after flag.Parse().
//
// Stop must be called for graceful shutdown.
func Init ( ) {
if len ( * remoteWriteURLs ) == 0 {
logger . Panicf ( "FATAL: at least one `-remoteWrite.url` must be set" )
}
if ! * showRemoteWriteURL {
// remoteWrite.url can contain authentication codes, so hide it at `/metrics` output.
httpserver . RegisterSecretFlag ( "remoteWrite.url" )
}
2020-03-03 12:08:17 +01:00
initRelabelGlobal ( )
2020-02-23 12:35:47 +01:00
maxInmemoryBlocks := memory . Allowed ( ) / len ( * remoteWriteURLs ) / maxRowsPerBlock / 100
if maxInmemoryBlocks > 200 {
// There is no much sense in keeping higher number of blocks in memory,
// since this means that the producer outperforms consumer and the queue
// will continue growing. It is better storing the queue to file.
maxInmemoryBlocks = 200
}
if maxInmemoryBlocks < 2 {
maxInmemoryBlocks = 2
}
for i , remoteWriteURL := range * remoteWriteURLs {
2020-03-03 12:08:17 +01:00
relabelConfigPath := ""
if i < len ( * relabelConfigPaths ) {
relabelConfigPath = ( * relabelConfigPaths ) [ i ]
}
2020-02-23 12:35:47 +01:00
urlLabelValue := fmt . Sprintf ( "secret-url-%d" , i + 1 )
if * showRemoteWriteURL {
urlLabelValue = remoteWriteURL
}
2020-05-06 15:51:32 +02:00
rwctx := newRemoteWriteCtx ( i , remoteWriteURL , relabelConfigPath , maxInmemoryBlocks , urlLabelValue )
2020-03-03 12:08:17 +01:00
rwctxs = append ( rwctxs , rwctx )
2020-02-23 12:35:47 +01:00
}
}
// Stop stops remotewrite.
//
// It is expected that nobody calls Push during and after the call to this func.
func Stop ( ) {
2020-03-03 12:08:17 +01:00
for _ , rwctx := range rwctxs {
rwctx . MustStop ( )
2020-02-23 12:35:47 +01:00
}
2020-03-03 12:08:17 +01:00
rwctxs = nil
2020-02-23 12:35:47 +01:00
}
// Push sends wr to remote storage systems set via `-remoteWrite.url`.
//
// Each timeseries in wr.Timeseries must contain one sample.
func Push ( wr * prompbmarshal . WriteRequest ) {
2020-03-03 12:08:17 +01:00
var rctx * relabelCtx
2020-03-06 18:26:17 +01:00
if len ( prcsGlobal ) > 0 || len ( labelsGlobal ) > 0 {
2020-03-03 12:08:17 +01:00
rctx = getRelabelCtx ( )
}
2020-02-28 17:57:45 +01:00
tss := wr . Timeseries
for len ( tss ) > 0 {
// Process big tss in smaller blocks in order to reduce maxmimum memory usage
tssBlock := tss
if len ( tssBlock ) > maxRowsPerBlock {
tssBlock = tss [ : maxRowsPerBlock ]
tss = tss [ maxRowsPerBlock : ]
2020-02-28 19:03:38 +01:00
} else {
tss = nil
2020-02-28 17:57:45 +01:00
}
2020-03-03 12:08:17 +01:00
if rctx != nil {
tssBlockLen := len ( tssBlock )
tssBlock = rctx . applyRelabeling ( tssBlock , labelsGlobal , prcsGlobal )
globalRelabelMetricsDropped . Add ( tssBlockLen - len ( tssBlock ) )
}
for _ , rwctx := range rwctxs {
rwctx . Push ( tssBlock )
}
2020-03-03 14:00:52 +01:00
if rctx != nil {
rctx . reset ( )
}
2020-02-28 17:57:45 +01:00
}
2020-03-03 12:08:17 +01:00
if rctx != nil {
putRelabelCtx ( rctx )
}
2020-02-23 12:35:47 +01:00
}
2020-03-03 12:08:17 +01:00
var globalRelabelMetricsDropped = metrics . NewCounter ( "vmagent_remotewrite_global_relabel_metrics_dropped_total" )
type remoteWriteCtx struct {
fq * persistentqueue . FastQueue
c * client
prcs [ ] promrelabel . ParsedRelabelConfig
pss [ ] * pendingSeries
pssNextIdx uint64
relabelMetricsDropped * metrics . Counter
}
2020-05-06 15:51:32 +02:00
func newRemoteWriteCtx ( argIdx int , remoteWriteURL , relabelConfigPath string , maxInmemoryBlocks int , urlLabelValue string ) * remoteWriteCtx {
2020-03-03 12:08:17 +01:00
h := xxhash . Sum64 ( [ ] byte ( remoteWriteURL ) )
path := fmt . Sprintf ( "%s/persistent-queue/%016X" , * tmpDataPath , h )
2020-03-03 18:48:46 +01:00
fq := persistentqueue . MustOpenFastQueue ( path , remoteWriteURL , maxInmemoryBlocks , * maxPendingBytesPerURL )
_ = metrics . GetOrCreateGauge ( fmt . Sprintf ( ` vmagent_remotewrite_pending_data_bytes { path=%q, url=%q} ` , path , urlLabelValue ) , func ( ) float64 {
2020-03-03 12:08:17 +01:00
return float64 ( fq . GetPendingBytes ( ) )
} )
2020-03-03 18:48:46 +01:00
_ = metrics . GetOrCreateGauge ( fmt . Sprintf ( ` vmagent_remotewrite_pending_inmemory_blocks { path=%q, url=%q} ` , path , urlLabelValue ) , func ( ) float64 {
2020-03-03 12:08:17 +01:00
return float64 ( fq . GetInmemoryQueueLen ( ) )
} )
2020-05-06 15:51:32 +02:00
c := newClient ( argIdx , remoteWriteURL , urlLabelValue , fq , * queues )
2020-03-03 12:08:17 +01:00
var prcs [ ] promrelabel . ParsedRelabelConfig
if len ( relabelConfigPath ) > 0 {
var err error
prcs , err = promrelabel . LoadRelabelConfigs ( relabelConfigPath )
if err != nil {
logger . Panicf ( "FATAL: cannot load relabel configs from -remoteWrite.urlRelabelConfig=%q: %s" , relabelConfigPath , err )
}
}
pss := make ( [ ] * pendingSeries , * queues )
for i := range pss {
pss [ i ] = newPendingSeries ( fq . MustWriteBlock )
}
return & remoteWriteCtx {
fq : fq ,
c : c ,
prcs : prcs ,
pss : pss ,
2020-03-03 18:48:46 +01:00
relabelMetricsDropped : metrics . GetOrCreateCounter ( fmt . Sprintf ( ` vmagent_remotewrite_relabel_metrics_dropped_total { path=%q, url=%q} ` , path , urlLabelValue ) ) ,
2020-02-23 12:35:47 +01:00
}
}
2020-03-03 12:08:17 +01:00
func ( rwctx * remoteWriteCtx ) MustStop ( ) {
for _ , ps := range rwctx . pss {
ps . MustStop ( )
}
rwctx . pss = nil
rwctx . fq . MustClose ( )
rwctx . fq = nil
rwctx . prcs = nil
rwctx . c . MustStop ( )
rwctx . c = nil
rwctx . relabelMetricsDropped = nil
}
2020-02-23 12:35:47 +01:00
2020-03-03 12:08:17 +01:00
func ( rwctx * remoteWriteCtx ) Push ( tss [ ] prompbmarshal . TimeSeries ) {
var rctx * relabelCtx
if len ( rwctx . prcs ) > 0 {
rctx = getRelabelCtx ( )
tssLen := len ( tss )
tss = rctx . applyRelabeling ( tss , nil , rwctx . prcs )
rwctx . relabelMetricsDropped . Add ( tssLen - len ( tss ) )
}
pss := rwctx . pss
idx := atomic . AddUint64 ( & rwctx . pssNextIdx , 1 ) % uint64 ( len ( pss ) )
pss [ idx ] . Push ( tss )
if rctx != nil {
putRelabelCtx ( rctx )
}
}