mirror of
https://github.com/VictoriaMetrics/VictoriaMetrics.git
synced 2026-08-24 20:27:48 +03:00
Compare commits
5 Commits
docs-updat
...
fixed-rule
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f881315f3f | ||
|
|
b7ef6cc3ed | ||
|
|
420f18b219 | ||
|
|
364b3b6823 | ||
|
|
400f01a114 |
@@ -13,6 +13,7 @@ import (
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/app/vmselect"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/app/vmselect/promql"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/app/vmstorage"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/appmetrics"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/buildinfo"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/cgroup"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/envflag"
|
||||
@@ -34,9 +35,10 @@ var (
|
||||
"This can be changed with -promscrape.config.strictParse=false command-line flag")
|
||||
maxIngestionRate = flag.Int("maxIngestionRate", 0, "The maximum number of samples vmsingle can receive per second. Data ingestion is paused when the limit is exceeded. "+
|
||||
"By default there are no limits on samples ingestion rate.")
|
||||
vmselectMaxConcurrentRequests = flag.Int("search.maxConcurrentRequests", getDefaultMaxConcurrentRequests(), "The maximum number of concurrent search requests. "+
|
||||
"It shouldn't be high, since a single request can saturate all the CPU cores, while many concurrently executed requests may require high amounts of memory. "+
|
||||
"See also -search.maxQueueDuration and -search.maxMemoryPerQuery")
|
||||
vmselectMaxConcurrentRequests = flagutil.NewIntWithDynamicDefault("search.maxConcurrentRequests", getDefaultMaxConcurrentRequests(), "vmselect.getDefaultMaxConcurrentRequests()",
|
||||
"The maximum number of concurrent search requests. "+
|
||||
"It shouldn't be high, since a single request can saturate all the CPU cores, while many concurrently executed requests may require high amounts of memory. "+
|
||||
"See also -search.maxQueueDuration and -search.maxMemoryPerQuery")
|
||||
vmselectMaxQueueDuration = flag.Duration("search.maxQueueDuration", 10*time.Second, "The maximum time the request waits for execution when -search.maxConcurrentRequests "+
|
||||
"limit is reached; see also -search.maxQueryDuration")
|
||||
)
|
||||
@@ -90,7 +92,9 @@ func main() {
|
||||
}
|
||||
logger.Infof("starting VictoriaMetrics at %q...", listenAddrs)
|
||||
startTime := time.Now()
|
||||
|
||||
vmstorage.Init(*vmselectMaxConcurrentRequests, *vmselectMaxQueueDuration, promql.ResetRollupResultCacheIfNeeded)
|
||||
appmetrics.MustCreateUncleanShutdownMarker(vmstorage.DataPath())
|
||||
vmselect.Init(*vmselectMaxConcurrentRequests, *vmselectMaxQueueDuration)
|
||||
vminsertcommon.StartIngestionRateLimiter(*maxIngestionRate)
|
||||
vminsert.Init()
|
||||
@@ -120,6 +124,7 @@ func main() {
|
||||
|
||||
vmstorage.Stop()
|
||||
vmselect.Stop()
|
||||
appmetrics.MustRemoveUncleanShutdownMarker(vmstorage.DataPath())
|
||||
|
||||
logger.Infof("the VictoriaMetrics has been stopped in %.3f seconds", time.Since(startTime).Seconds())
|
||||
}
|
||||
|
||||
@@ -15,6 +15,7 @@ import (
|
||||
"github.com/VictoriaMetrics/metrics"
|
||||
"github.com/cespare/xxhash/v2"
|
||||
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/appmetrics"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/auth"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/bloomfilter"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/cgroup"
|
||||
@@ -62,9 +63,10 @@ var (
|
||||
"See also -remoteWrite.maxDiskUsagePerURL and -remoteWrite.disableOnDiskQueue")
|
||||
keepDanglingQueues = flag.Bool("remoteWrite.keepDanglingQueues", false, "Keep persistent queues contents at -remoteWrite.tmpDataPath in case there are no matching -remoteWrite.url. "+
|
||||
"Useful when -remoteWrite.url is changed temporarily and persistent queue files will be needed later on.")
|
||||
queues = flagutil.NewArrayInt("remoteWrite.queues", cgroup.AvailableCPUs()*2, "The number of concurrent queues to each -remoteWrite.url. Set more queues if default number of queues "+
|
||||
"isn't enough for sending high volume of collected data to remote storage. "+
|
||||
"Default value depends on the number of available CPU cores. It should work fine in most cases since it minimizes resource usage")
|
||||
queues = flagutil.NewArrayIntWithDynamicDefault("remoteWrite.queues", cgroup.AvailableCPUs()*2, "2*cgroup.AvailableCPUs()",
|
||||
"The number of concurrent queues to each -remoteWrite.url. Set more queues if default number of queues "+
|
||||
"isn't enough for sending high volume of collected data to remote storage. "+
|
||||
"Default value depends on the number of available CPU cores. It should work fine in most cases since it minimizes resource usage")
|
||||
inmemoryQueues = flagutil.NewArrayInt("remoteWrite.inmemoryQueues", 0, "The number of additional workers per each -remoteWrite.url, which send only recently ingested data from the in-memory queue, "+
|
||||
"while the file-based queue at -remoteWrite.tmpDataPath is drained by workers configured via -remoteWrite.queues. "+
|
||||
"This reduces delivery lag for fresh samples when the file-based queue contains a backlog accumulated during remote storage outages.")
|
||||
@@ -233,6 +235,7 @@ func Init() {
|
||||
initStreamAggrConfigGlobal()
|
||||
|
||||
initRemoteWriteCtxs(*remoteWriteURLs)
|
||||
appmetrics.MustCreateUncleanShutdownMarker(*tmpDataPath)
|
||||
|
||||
disableOnDiskQueues := []bool(*disableOnDiskQueue)
|
||||
disableOnDiskQueueAny = slices.Contains(disableOnDiskQueues, true)
|
||||
@@ -391,6 +394,8 @@ func Stop() {
|
||||
if sl := dailySeriesLimiter; sl != nil {
|
||||
sl.MustStop()
|
||||
}
|
||||
|
||||
appmetrics.MustRemoveUncleanShutdownMarker(*tmpDataPath)
|
||||
}
|
||||
|
||||
// PushDropSamplesOnFailure pushes wr to the configured remote storage systems set via -remoteWrite.url
|
||||
|
||||
@@ -36,9 +36,10 @@ var (
|
||||
idleConnectionTimeout = flag.Duration("remoteWrite.idleConnTimeout", 50*time.Second, `Defines a duration for idle (keep-alive connections) to exist. Consider settings this value less to the value of "-http.idleConnTimeout". It must prevent possible "write: broken pipe" and "read: connection reset by peer" errors.`)
|
||||
maxIdleConnections = flag.Int("remoteWrite.maxIdleConnections", 100, `Defines the number of idle (keep-alive connections) to -remoteWrite.url for the vmalert-tool debug writer, which sends every series in a separate request. Too low a value may result in a high number of sockets in TIME_WAIT state.`)
|
||||
|
||||
maxQueueSize = flag.Int("remoteWrite.maxQueueSize", defaultMaxQueueSize, "Defines the max number of pending datapoints to remote write endpoint")
|
||||
maxBatchSize = flag.Int("remoteWrite.maxBatchSize", defaultMaxBatchSize, "Defines max number of timeseries to be flushed at once")
|
||||
concurrency = flag.Int("remoteWrite.concurrency", defaultConcurrency, "Defines number of writers for concurrent writing into remote write endpoint. Default value depends on the number of available CPU cores.")
|
||||
maxQueueSize = flag.Int("remoteWrite.maxQueueSize", defaultMaxQueueSize, "Defines the max number of pending datapoints to remote write endpoint")
|
||||
maxBatchSize = flag.Int("remoteWrite.maxBatchSize", defaultMaxBatchSize, "Defines max number of timeseries to be flushed at once")
|
||||
concurrency = flagutil.NewIntWithDynamicDefault("remoteWrite.concurrency", defaultConcurrency, "2*cgroup.AvailableCPUs()",
|
||||
"Defines number of writers for concurrent writing into remote write endpoint. Default value depends on the number of available CPU cores.")
|
||||
flushInterval = flag.Duration("remoteWrite.flushInterval", defaultFlushInterval, "Defines interval of flushes to remote write endpoint")
|
||||
|
||||
tlsInsecureSkipVerify = flag.Bool("remoteWrite.tlsInsecureSkipVerify", false, "Whether to skip tls verification when connecting to -remoteWrite.url")
|
||||
|
||||
@@ -20,6 +20,7 @@ import (
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/bytesutil"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/cgroup"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/fasttime"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/flagutil"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/querytracer"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/storage"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/storage/metricnamestats"
|
||||
@@ -30,11 +31,12 @@ var (
|
||||
maxSamplesPerSeries = flag.Int("search.maxSamplesPerSeries", 30e6, "The maximum number of raw samples a single query can scan per each time series. This option allows limiting memory usage")
|
||||
maxSamplesPerQuery = flag.Int("search.maxSamplesPerQuery", 1e9, "The maximum number of raw samples a single query can process across all time series. "+
|
||||
"This protects from heavy queries, which select unexpectedly high number of raw samples. See also -search.maxSamplesPerSeries")
|
||||
maxWorkersPerQuery = flag.Int("search.maxWorkersPerQuery", defaultMaxWorkersPerQuery, "The maximum number of CPU cores a single query can use. "+
|
||||
"The default value should work good for most cases. "+
|
||||
"The flag can be set to lower values for improving performance of big number of concurrently executed queries. "+
|
||||
"The flag can be set to bigger values for improving performance of heavy queries, which scan big number of time series (>10K) and/or big number of samples (>100M). "+
|
||||
"There is no sense in setting this flag to values bigger than the number of CPU cores available on the system")
|
||||
maxWorkersPerQuery = flagutil.NewIntWithDynamicDefault("search.maxWorkersPerQuery", defaultMaxWorkersPerQuery, "netstorage.defaultMaxWorkersPerQuery()",
|
||||
"The maximum number of CPU cores a single query can use. "+
|
||||
"The default value should work good for most cases. "+
|
||||
"The flag can be set to lower values for improving performance of big number of concurrently executed queries. "+
|
||||
"The flag can be set to bigger values for improving performance of heavy queries, which scan big number of time series (>10K) and/or big number of samples (>100M). "+
|
||||
"There is no sense in setting this flag to values bigger than the number of CPU cores available on the system")
|
||||
)
|
||||
|
||||
// Result is a single timeseries result.
|
||||
|
||||
@@ -62,6 +62,9 @@ type tmpBlocksFile struct {
|
||||
r *fs.ReaderAt
|
||||
|
||||
offset uint64
|
||||
|
||||
// err stores the first error occurred while writing the temporary blocks file.
|
||||
err error
|
||||
}
|
||||
|
||||
func getTmpBlocksFile() *tmpBlocksFile {
|
||||
@@ -82,6 +85,7 @@ func putTmpBlocksFile(tbf *tmpBlocksFile) {
|
||||
tbf.f = nil
|
||||
tbf.r = nil
|
||||
tbf.offset = 0
|
||||
tbf.err = nil
|
||||
tmpBlocksFilePool.Put(tbf)
|
||||
}
|
||||
|
||||
@@ -109,6 +113,10 @@ var (
|
||||
// and this must be handled.
|
||||
func (tbf *tmpBlocksFile) WriteBlockRefData(b []byte) (tmpBlockAddr, error) {
|
||||
var addr tmpBlockAddr
|
||||
if tbf.err != nil {
|
||||
// Do not write anything to the tbf after the first failed write
|
||||
return addr, tbf.err
|
||||
}
|
||||
addr.offset = tbf.offset
|
||||
addr.size = len(b)
|
||||
tbf.offset += uint64(addr.size)
|
||||
@@ -122,7 +130,8 @@ func (tbf *tmpBlocksFile) WriteBlockRefData(b []byte) (tmpBlockAddr, error) {
|
||||
if tbf.f == nil {
|
||||
f, err := os.CreateTemp(tmpBlocksDir, "")
|
||||
if err != nil {
|
||||
return addr, err
|
||||
tbf.err = fmt.Errorf("cannot create temporary blocks file at %q: %w", tmpBlocksDir, err)
|
||||
return addr, tbf.err
|
||||
}
|
||||
tbf.f = f
|
||||
tmpBlocksFilesCreated.Inc()
|
||||
@@ -130,7 +139,9 @@ func (tbf *tmpBlocksFile) WriteBlockRefData(b []byte) (tmpBlockAddr, error) {
|
||||
_, err := tbf.f.Write(tbf.buf)
|
||||
tbf.buf = append(tbf.buf[:0], b...)
|
||||
if err != nil {
|
||||
return addr, fmt.Errorf("cannot write block to %q: %w", tbf.f.Name(), err)
|
||||
// The blocks buffered at tbf.buf could be partially lost, mark the tbf as unusable.
|
||||
tbf.err = fmt.Errorf("cannot write block to %q: %w", tbf.f.Name(), err)
|
||||
return addr, tbf.err
|
||||
}
|
||||
return addr, nil
|
||||
}
|
||||
@@ -141,12 +152,16 @@ func (tbf *tmpBlocksFile) Len() uint64 {
|
||||
}
|
||||
|
||||
func (tbf *tmpBlocksFile) Finalize() error {
|
||||
if tbf.err != nil {
|
||||
return tbf.err
|
||||
}
|
||||
if tbf.f == nil {
|
||||
return nil
|
||||
}
|
||||
fname := tbf.f.Name()
|
||||
if _, err := tbf.f.Write(tbf.buf); err != nil {
|
||||
return fmt.Errorf("cannot write the remaining %d bytes to %q: %w", len(tbf.buf), fname, err)
|
||||
tbf.err = fmt.Errorf("cannot write the remaining %d bytes to %q: %w", len(tbf.buf), fname, err)
|
||||
return tbf.err
|
||||
}
|
||||
tbf.buf = tbf.buf[:0]
|
||||
r := fs.NewReaderAt(tbf.f)
|
||||
@@ -166,6 +181,10 @@ func (tbf *tmpBlocksFile) Finalize() error {
|
||||
}
|
||||
|
||||
func (tbf *tmpBlocksFile) MustReadBlockRefAt(partRef storage.PartRef, addr tmpBlockAddr) storage.BlockRef {
|
||||
if tbf.err != nil {
|
||||
// This should never happen, since Finalize() already returns the error for such a tbf.
|
||||
logger.Panicf("BUG: cannot read block at %s from the temporary blocks file with the failed write: %s", addr, tbf.err)
|
||||
}
|
||||
var buf []byte
|
||||
if tbf.r == nil {
|
||||
buf = tbf.buf[addr.offset : addr.offset+uint64(addr.size)]
|
||||
|
||||
@@ -16,6 +16,19 @@ groups:
|
||||
Job {{ $labels.job }} (instance {{ $labels.instance }}) has restarted more than twice in the last 15 minutes.
|
||||
It might be crashlooping.
|
||||
|
||||
- alert: UncleanShutdown
|
||||
expr: vm_app_prev_shutdown_unclean == 1 and time() - vm_app_start_timestamp < 600
|
||||
labels:
|
||||
severity: warning
|
||||
annotations:
|
||||
summary: "{{ $labels.job }} on instance {{ $labels.instance }} started after an unclean shutdown"
|
||||
description: |
|
||||
The previous process run didn't shut down cleanly. Check the logs for OOM, SIGKILL,
|
||||
a host failure, or another unexpected termination. In Kubernetes, a pod may be forcefully
|
||||
killed with SIGKILL if the shutdown takes longer than terminationGracePeriodSeconds.
|
||||
This alert stops firing 10 minutes after startup.
|
||||
See https://github.com/VictoriaMetrics/VictoriaMetrics/issues/8443 for more details.
|
||||
|
||||
- alert: ServiceDown
|
||||
expr: up{job=~".*(victoriametrics|vmselect|vminsert|vmstorage|vmagent|vmalert|vmsingle|vmalertmanager|vmauth).*"} == 0
|
||||
for: 2m
|
||||
|
||||
@@ -126,10 +126,10 @@ groups:
|
||||
(
|
||||
vmalert_alerting_rules_last_evaluation_samples
|
||||
> on(group,file) group_left()
|
||||
(vmalert_group_rule_results_limit * 0.9)
|
||||
min by (group,file) (vmalert_group_rule_results_limit * 0.9)
|
||||
)
|
||||
and on(group,file)
|
||||
(vmalert_group_rule_results_limit > 0)
|
||||
(min by (group,file) (vmalert_group_rule_results_limit) > 0)
|
||||
for: 5m
|
||||
labels:
|
||||
severity: warning
|
||||
@@ -144,10 +144,10 @@ groups:
|
||||
(
|
||||
vmalert_recording_rules_last_evaluation_samples
|
||||
> on(group,file) group_left()
|
||||
(vmalert_group_rule_results_limit * 0.9)
|
||||
min by (group,file) (vmalert_group_rule_results_limit * 0.9)
|
||||
)
|
||||
and on(group,file)
|
||||
(vmalert_group_rule_results_limit > 0)
|
||||
(min by (group,file) (vmalert_group_rule_results_limit) > 0)
|
||||
for: 5m
|
||||
labels:
|
||||
severity: warning
|
||||
|
||||
@@ -90,12 +90,9 @@ endif
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/victoria_metrics_common_flags.md
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/victoria_metrics_enterprise_flags.md
|
||||
|
||||
# adjust flags with dynamic default values
|
||||
# remove after https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9680 implemented
|
||||
sed -i '/The maximum number of concurrent insert requests/ s/(default [0-9]\+)/(default 2*cgroup.AvailableCPUs())/' docs/victoriametrics/victoria_metrics_common_flags.md
|
||||
sed -i '/The maximum number of concurrent search requests\./ s/(default [0-9]\+)/(default vmselect.getDefaultMaxConcurrentRequests())/' docs/victoriametrics/victoria_metrics_common_flags.md
|
||||
sed -i '/The maximum number of CPU cores a single query can use\./ s/(default [0-9]\+)/(default netstorage.defaultMaxWorkersPerQuery())/' docs/victoriametrics/victoria_metrics_common_flags.md
|
||||
sed -i '/The maximum number of concurrent goroutines to work with files;/ s/(default [0-9]\+)/(default fsutil.getDefaultConcurrency())/' docs/victoriametrics/victoria_metrics_common_flags.md
|
||||
# hide the machine-specific value of dynamic defaults, keeping the formula.
|
||||
# the flagutil.New*WithDynamicDefault constructors print them as "(default <value> = <formula>)".
|
||||
sed -i 's/(default [0-9]\+ = \(.*\))$$/(default \1)/' docs/victoriametrics/victoria_metrics_common_flags.md
|
||||
|
||||
docs-update-vmauth-flags:
|
||||
ifndef TAG
|
||||
@@ -119,9 +116,9 @@ endif
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/vmauth_common_flags.md
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/vmauth_enterprise_flags.md
|
||||
|
||||
# adjust flags with dynamic default values
|
||||
# remove after https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9680 implemented
|
||||
sed -i '/The maximum number of concurrent goroutines to work with files;/ s/(default [0-9]\+)/(default fsutil.getDefaultConcurrency())/' docs/victoriametrics/vmauth_common_flags.md
|
||||
# hide the machine-specific value of dynamic defaults, keeping the formula.
|
||||
# the flagutil.New*WithDynamicDefault constructors print them as "(default <value> = <formula>)".
|
||||
sed -i 's/(default [0-9]\+ = \(.*\))$$/(default \1)/' docs/victoriametrics/vmauth_common_flags.md
|
||||
|
||||
docs-update-vmagent-flags:
|
||||
ifndef TAG
|
||||
@@ -145,11 +142,9 @@ endif
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/vmagent_common_flags.md
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/vmagent_enterprise_flags.md
|
||||
|
||||
# adjust flags with dynamic default values
|
||||
# remove after https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9680 implemented
|
||||
sed -i '/The maximum number of concurrent insert requests/ s/(default [0-9]\+)/(default 2*cgroup.AvailableCPUs())/' docs/victoriametrics/vmagent_common_flags.md
|
||||
sed -i '/The number of concurrent queues to each -remoteWrite.url./ s/(default [0-9]\+)/(default 2*cgroup.AvailableCPUs())/' docs/victoriametrics/vmagent_common_flags.md
|
||||
sed -i '/The maximum number of concurrent goroutines to work with files;/ s/(default [0-9]\+)/(default fsutil.getDefaultConcurrency())/' docs/victoriametrics/vmagent_common_flags.md
|
||||
# hide the machine-specific value of dynamic defaults, keeping the formula.
|
||||
# the flagutil.New*WithDynamicDefault constructors print them as "(default <value> = <formula>)".
|
||||
sed -i 's/(default [0-9]\+ = \(.*\))$$/(default \1)/' docs/victoriametrics/vmagent_common_flags.md
|
||||
|
||||
docs-update-vmalert-flags:
|
||||
ifndef TAG
|
||||
@@ -173,10 +168,9 @@ endif
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/vmalert_common_flags.md
|
||||
sed -i 's/\t/ /g' docs/victoriametrics/vmalert_enterprise_flags.md
|
||||
|
||||
# adjust flags with dynamic default values
|
||||
# remove after https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9680 implemented
|
||||
sed -i '/Defines number of writers for concurrent writing into remote write endpoint./ s/(default [0-9]\+)/(default 2*cgroup.AvailableCPUs())/' docs/victoriametrics/vmalert_common_flags.md
|
||||
sed -i '/The maximum number of concurrent goroutines to work with files;/ s/(default [0-9]\+)/(default fsutil.getDefaultConcurrency())/' docs/victoriametrics/vmalert_common_flags.md
|
||||
# hide the machine-specific value of dynamic defaults, keeping the formula.
|
||||
# the flagutil.New*WithDynamicDefault constructors print them as "(default <value> = <formula>)".
|
||||
sed -i 's/(default [0-9]\+ = \(.*\))$$/(default \1)/' docs/victoriametrics/vmalert_common_flags.md
|
||||
|
||||
docs-update-vmselect-flags:
|
||||
ifndef TAG
|
||||
|
||||
@@ -300,26 +300,26 @@ See [VMDistributed](https://docs.victoriametrics.com/operator/resources/vmdistri
|
||||
|
||||
## Cluster setup
|
||||
|
||||
A minimal cluster setup consists of the following components:
|
||||
A minimal cluster must contain the following nodes:
|
||||
|
||||
- a single `vmstorage` node with `-retentionPeriod` and `-storageDataPath` flags
|
||||
- a single `vminsert` node with `-storageNode=<vmstorage_host>`
|
||||
- a single `vmselect` node with `-storageNode=<vmstorage_host>`
|
||||
|
||||
> The [Enterprise version of VictoriaMetrics](https://docs.victoriametrics.com/victoriametrics/enterprise/) supports automatic discovery and updating of `vmstorage` nodes.
|
||||
> See [automatic vmstorage discovery](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#automatic-vmstorage-discovery) for details.
|
||||
[Enterprise version of VictoriaMetrics](https://docs.victoriametrics.com/victoriametrics/enterprise/) supports automatic discovering and updating of `vmstorage` nodes.
|
||||
See [these docs](#automatic-vmstorage-discovery) for details.
|
||||
|
||||
Prefer running many small `vmstorage` nodes over a few big `vmstorage` nodes. For example, prefer running at least 10 `vmstorage` nodes for better [load distribution and availability](#cluster-availability). If you need fewer `vmstorage` nodes, consider using the [single-node version](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/) of VictoriaMetrics instead. See more on [choosing between single-node and cluster versions](https://docs.victoriametrics.com/victoriametrics/faq/#which-victoriametrics-type-is-recommended-for-use-in-production---single-node-or-cluster).
|
||||
It is recommended to run at least two nodes for each service for high availability purposes. In this case the cluster continues working when a single node is temporarily unavailable and the remaining nodes can handle the increased workload. The node may be temporarily unavailable when the underlying hardware breaks, during software upgrades, migration or other maintenance tasks.
|
||||
|
||||
If you run multiple nodes of `vminsert` or `vmselect`, use an HTTP load balancer such as [vmauth](https://docs.victoriametrics.com/victoriametrics/vmauth/)
|
||||
or `nginx` in front of them. It must contain the following routing configs according to [the URL format](#url-format):
|
||||
It is preferred to run many small `vmstorage` nodes over a few big `vmstorage` nodes, since this reduces the workload increase on the remaining `vmstorage` nodes when some of `vmstorage` nodes become temporarily unavailable.
|
||||
|
||||
An http load balancer such as [vmauth](https://docs.victoriametrics.com/victoriametrics/vmauth/) or `nginx` must be put in front of `vminsert` and `vmselect` nodes.
|
||||
It must contain the following routing configs according to [the url format](#url-format):
|
||||
|
||||
- requests starting with `/insert` must be routed to port `8480` on `vminsert` nodes.
|
||||
- requests starting with `/select` must be routed to port `8481` on `vmselect` nodes.
|
||||
|
||||
> Ports may be altered by setting `-httpListenAddr` on the corresponding nodes.
|
||||
|
||||
See an example of the [vmauth configuration for a VictoriaMetrics cluster](https://docs.victoriametrics.com/vmauth/index.html#load-balancer-for-victoriametrics-cluster).
|
||||
Ports may be altered by setting `-httpListenAddr` on the corresponding nodes.
|
||||
|
||||
It is recommended setting up [monitoring](#monitoring) for the cluster.
|
||||
|
||||
@@ -346,32 +346,35 @@ VictoriaMetrics cluster remains available if the following conditions are met:
|
||||
- HTTP load balancer must stop routing requests to unavailable `vminsert` and `vmselect` nodes
|
||||
([vmauth](https://docs.victoriametrics.com/victoriametrics/vmauth/) stops routing requests to unavailable nodes).
|
||||
|
||||
- At least a single `vminsert` node must remain available in the cluster for processing the ingestion workload.
|
||||
- At least a single `vminsert` node must remain available in the cluster for processing data ingestion workload.
|
||||
The remaining active `vminsert` nodes must have enough compute capacity (CPU, RAM, network bandwidth)
|
||||
for handling the ingestion workload. Otherwise, ingestion could stop or be delayed.
|
||||
for handling the current data ingestion workload.
|
||||
If the remaining active `vminsert` nodes have no enough resources for processing the data ingestion workload,
|
||||
then arbitrary delays may occur during data ingestion.
|
||||
See [capacity planning](#capacity-planning) and [cluster resizing](#cluster-resizing-and-scalability) docs for more details.
|
||||
|
||||
- At least a single `vmselect` node must remain available in the cluster for processing query workload.
|
||||
The remaining active `vmselect` nodes must have enough compute capacity (CPU, RAM, network bandwidth, disk IO)
|
||||
for handling the query workload.
|
||||
If the remaining active `vmselect` nodes do not have enough resources for processing the query workload,
|
||||
then arbitrary query failures and latency increases may occur during query processing.
|
||||
for handling the current query workload.
|
||||
If the remaining active `vmselect` nodes have no enough resources for processing query workload,
|
||||
then arbitrary failures and delays may occur during query processing.
|
||||
See [capacity planning](#capacity-planning) and [cluster resizing](#cluster-resizing-and-scalability) docs for more details.
|
||||
|
||||
- At least a single `vmstorage` node must remain available in the cluster for accepting newly ingested data
|
||||
and for processing incoming read queries. The remaining active `vmstorage` nodes must have enough compute capacity
|
||||
(CPU, RAM, network bandwidth, disk IO, free disk space) for handling the workload.
|
||||
If the remaining active `vmstorage` nodes do not have enough resources for processing the workload,
|
||||
then arbitrary failures and delays may occur during data ingestion and query processing.
|
||||
and for processing incoming queries. The remaining active `vmstorage` nodes must have enough compute capacity
|
||||
(CPU, RAM, network bandwidth, disk IO, free disk space) for handling the current workload.
|
||||
If the remaining active `vmstorage` nodes have no enough resources for processing query workload,
|
||||
then arbitrary failures and delay may occur during data ingestion and query processing.
|
||||
See [capacity planning](#capacity-planning) and [cluster resizing](#cluster-resizing-and-scalability) docs for more details.
|
||||
|
||||
The cluster works in the following way when some of `vmstorage` nodes are unavailable:
|
||||
|
||||
- `vminsert` [re-routes](https://victoriametrics.com/blog/vminsert-how-it-works/#31-rerouting) newly ingested data from
|
||||
unavailable `vmstorage` nodes to remaining healthy `vmstorage` nodes. This guarantees that the newly ingested data is
|
||||
properly saved if the healthy `vmstorage` nodes have enough CPU, RAM, disk I/O, and network bandwidth for processing
|
||||
the increased data ingestion workload. During re-routing, healthy `vmstorage` nodes will experience higher resource usage
|
||||
and an increase in the number of [active time series](https://docs.victoriametrics.com/victoriametrics/faq/#what-is-an-active-time-series).
|
||||
- `vminsert` re-routes newly ingested data from unavailable `vmstorage` nodes to remaining healthy `vmstorage` nodes.
|
||||
This guarantees that the newly ingested data is properly saved if the healthy `vmstorage` nodes have enough CPU, RAM, disk IO and network bandwidth
|
||||
for processing the increased data ingestion workload.
|
||||
`vminsert` spreads evenly the additional data among the healthy `vmstorage` nodes in order to spread evenly
|
||||
the increased load on these nodes. During re-routing, healthy `vmstorage` nodes will experience higher resource usage
|
||||
and increase in number of [active time series](https://docs.victoriametrics.com/victoriametrics/faq/#what-is-an-active-time-series).
|
||||
|
||||
- `vmselect` continues serving queries if at least a single `vmstorage` nodes is available.
|
||||
It marks responses as partial for queries served from the remaining healthy `vmstorage` nodes,
|
||||
|
||||
@@ -442,29 +442,29 @@ Both [single-node VictoriaMetrics](https://docs.victoriametrics.com/victoriametr
|
||||
|
||||
See [Scalability limits of VictoriaMetrics](https://docs.victoriametrics.com/victoriametrics/faq/#what-are-scalability-limits-of-victoriametrics).
|
||||
|
||||
Benefits of using single-node VictoriaMetrics:
|
||||
Single-node VictoriaMetrics requires lower amounts of CPU and RAM for handling the same workload comparing
|
||||
to cluster version of VictoriaMetrics, since it doesn't need to pass the encoded data over the network
|
||||
between [cluster components](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#architecture-overview).
|
||||
|
||||
* Requires less CPU and RAM than the cluster version for the same workload because it doesn't need to transfer encoded data over the network between [cluster components](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#architecture-overview)
|
||||
The performance of a single-node VictoriaMetrics scales almost perfectly with the available CPU, RAM and disk IO resources on the host where it runs -
|
||||
see [this article](https://valyala.medium.com/measuring-vertical-scalability-for-time-series-databases-in-google-cloud-92550d78d8ae).
|
||||
|
||||
* Scales almost perfectly with the available CPU, RAM, and disk I/O resources on the host where it runs. See [Measuring vertical scalability for time series databases in Google Cloud](https://valyala.medium.com/measuring-vertical-scalability-for-time-series-databases-in-google-cloud-92550d78d8ae) to learn how VictoriaMetrics single-node scales vertically.
|
||||
|
||||
* It is easier to set up and operate compared to the cluster version of VictoriaMetrics
|
||||
|
||||
* Supports [high availability](https://docs.victoriametrics.com/Single-server-VictoriaMetrics/#high-availability)
|
||||
Single-node VictoriaMetrics is easier to setup and operate comparing to cluster version of VictoriaMetrics.
|
||||
|
||||
Given the facts above **it is recommended to use single-node VictoriaMetrics in the majority of cases**.
|
||||
|
||||
Cluster version of VictoriaMetrics may be preferred over single-node VictoriaMetrics in the following relatively rare cases:
|
||||
|
||||
* If [multitenancy support](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#multitenancy) is needed.
|
||||
Single-node VictoriaMetrics doesn't support multitenancy. Though it is possible to run multiple single-node VictoriaMetrics
|
||||
instances - one per each tenant. See more about [multitenancy in the single-node version](https://docs.victoriametrics.com/victoriametrics/#multi-tenancy). - and route incoming requests from a particular tenant to the needed VictoriaMetrics instance
|
||||
Multitenancy can also be achieved via metric labels (i.e., `{env="prod"}` or `{team="platform"}`) and [enforcing](https://docs.victoriametrics.com/victoriametrics/#prometheus-querying-api-enhancements)
|
||||
labels filters on reads via [vmauth](https://docs.victoriametrics.com/victoriametrics/vmauth/#enforcing-query-args).
|
||||
* If [multitenancy support](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#multitenancy) is needed,
|
||||
since single-node VictoriaMetrics doesn't support multitenancy. Though it is possible to run multiple single-node VictoriaMetrics
|
||||
instances - one per each tenant - and route incoming requests from particular tenant to the needed VictoriaMetrics instance
|
||||
via [vmauth](https://docs.victoriametrics.com/victoriametrics/vmauth/).
|
||||
|
||||
* If a single-node VictoriaMetrics cannot handle the current workload. For example, if you plan to ingest hundreds of millions of active time series at ingestion rates exceeding millions of samples per second, it is better to use a cluster version of VictoriaMetrics. Its capacity can [scale horizontally with the number of nodes in the cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#cluster-resizing-and-scalability).
|
||||
* If the current workload cannot be handled by a single-node VictoriaMetrics. For example, if you are going to ingest hundreds of millions of active time series
|
||||
at ingestion rates exceeding a million samples per second, then it is better to use cluster version of VictoriaMetrics,
|
||||
since its capacity can [scale horizontally with the number of nodes in the cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/#cluster-resizing-and-scalability).
|
||||
|
||||
[Don't choose VictoriaMetrics cluster unless you have to](https://victoriametrics.com/blog/dont-default-to-microservices-you-will-thank-us-later/).
|
||||
[Don't choose cluster unless you have to](https://victoriametrics.com/blog/dont-default-to-microservices-you-will-thank-us-later/).
|
||||
|
||||
## How to migrate data from single-node VictoriaMetrics to cluster version?
|
||||
|
||||
|
||||
@@ -26,11 +26,14 @@ See also [LTS releases](https://docs.victoriametrics.com/victoriametrics/lts-rel
|
||||
|
||||
## tip
|
||||
|
||||
* FEATURE: [vmagent](https://docs.victoriametrics.com/victoriametrics/vmagent/), [vmsingle](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/), `vmstorage` and `vmselect` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): expose the `vm_app_prev_shutdown_unclean` gauge. It is set to `1` when the previous process run didn't shut down cleanly. Added the `UncleanShutdown` [alerting rule](https://github.com/VictoriaMetrics/VictoriaMetrics/blob/master/deployment/docker/rules/alerts-health.yml), which fires for 10 minutes after an unclean shutdown is detected. See [#8443](https://github.com/VictoriaMetrics/VictoriaMetrics/issues/8443).
|
||||
* FEATURE: [vmui](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/#vmui): show the selected time zone UTC offset next to the date/time controls and allow opening time zone settings from it. See [#11332](https://github.com/VictoriaMetrics/VictoriaMetrics/pull/11332).
|
||||
* FEATURE: [vmsingle](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/), [vmagent](https://docs.victoriametrics.com/victoriametrics/vmagent/), [vmalert](https://docs.victoriametrics.com/victoriametrics/vmalert/), and `vmselect` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): show how the default value is calculated for command-line flags which derive it from the number of available CPU cores. For example, `-maxConcurrentInserts` now prints `(default 16 = 2*cgroup.AvailableCPUs())` in `-help` output instead of `(default 16)`. Updated flags: `-search.maxConcurrentRequests`, `-search.maxWorkersPerQuery`, `-fs.maxConcurrency`, `-remoteWrite.concurrency`, `-remoteWrite.queues`. See [#9680](https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9680). Thanks to @Vandit1604 for contribution.
|
||||
|
||||
* BUGFIX: [vmagent](https://docs.victoriametrics.com/victoriametrics/vmagent/) and `vminsert` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): fix infinite loop in the OpenTelemetry Firehose ingestion endpoint (`/opentelemetry/api/v1/push`) when receiving a malformed record with an incomplete varint in the `data` field. Previously this caused the goroutine to spin forever, permanently consuming CPU until the process was restarted.
|
||||
* BUGFIX: [vmalert-tool](https://docs.victoriametrics.com/victoriametrics/vmalert-tool/): reuse connections to `-remoteWrite.url` when writing the results of recording rules and alerts. Previously every series was sent over a new connection, which left a lot of sockets in `TIME_WAIT` state and could exhaust the ephemeral port range. The number of idle connections can be tuned via the new `-remoteWrite.maxIdleConnections` command-line flag. Thanks @evkuzin for contribution.
|
||||
* BUGFIX: [vmsingle](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/) and `vmselect` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): prevent process crash in `sort_by_label_numeric()` and `sort_by_label_numeric_desc()` when a label value contains a number with 309 or more digits. See [#11423](https://github.com/VictoriaMetrics/VictoriaMetrics/pull/11423).
|
||||
* BUGFIX: `vmselect` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): fail the query request directly when there is not enough disk space to store temporary search results. Previously, such queries could lead to vmselect crash. See [#4688](https://github.com/VictoriaMetrics/VictoriaMetrics/issues/4688).
|
||||
|
||||
## [v1.150.0](https://github.com/VictoriaMetrics/VictoriaMetrics/releases/tag/v1.150.0)
|
||||
|
||||
|
||||
@@ -65,6 +65,9 @@ func writePrometheusMetrics(w io.Writer) {
|
||||
// Export start time and uptime in seconds
|
||||
metrics.WriteGaugeUint64(w, "vm_app_start_timestamp", uint64(startTime.Unix()))
|
||||
metrics.WriteGaugeUint64(w, "vm_app_uptime_seconds", uint64(time.Since(startTime).Seconds()))
|
||||
if uncleanShutdownEnabled.Load() {
|
||||
metrics.WriteGaugeUint64(w, "vm_app_prev_shutdown_unclean", uncleanShutdown)
|
||||
}
|
||||
|
||||
// Export flags as metrics.
|
||||
isSetMap := make(map[string]bool)
|
||||
|
||||
@@ -13,13 +13,13 @@ type osInfo struct {
|
||||
release string
|
||||
}
|
||||
|
||||
var os osInfo
|
||||
var hostOS osInfo
|
||||
var initOSOnce sync.Once
|
||||
|
||||
func writeOSMetrics(w io.Writer) {
|
||||
initOSOnce.Do(initOS)
|
||||
|
||||
if os.name != "" {
|
||||
metrics.WriteGaugeUint64(w, fmt.Sprintf(`vm_os_info{os=%q, release=%q}`, os.name, os.release), 1)
|
||||
if hostOS.name != "" {
|
||||
metrics.WriteGaugeUint64(w, fmt.Sprintf(`vm_os_info{os=%q, release=%q}`, hostOS.name, hostOS.release), 1)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -8,7 +8,7 @@ import (
|
||||
)
|
||||
|
||||
func initOS() {
|
||||
os = osInfo{name: "darwin"}
|
||||
hostOS = osInfo{name: "darwin"}
|
||||
|
||||
out, err := exec.Command("sysctl", "-n", "kern.osrelease").Output()
|
||||
if err != nil {
|
||||
@@ -16,5 +16,5 @@ func initOS() {
|
||||
return
|
||||
}
|
||||
|
||||
os.release = strings.TrimSpace(string(out))
|
||||
hostOS.release = strings.TrimSpace(string(out))
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import (
|
||||
)
|
||||
|
||||
func initOS() {
|
||||
os = osInfo{name: "linux"}
|
||||
hostOS = osInfo{name: "linux"}
|
||||
|
||||
var uname syscall.Utsname
|
||||
if err := syscall.Uname(&uname); err != nil {
|
||||
@@ -22,5 +22,5 @@ func initOS() {
|
||||
}
|
||||
ur = append(ur, byte(v))
|
||||
}
|
||||
os.release = string(ur)
|
||||
hostOS.release = string(ur)
|
||||
}
|
||||
|
||||
@@ -8,12 +8,12 @@ import (
|
||||
)
|
||||
|
||||
func initOS() {
|
||||
os = osInfo{name: "windows"}
|
||||
hostOS = osInfo{name: "windows"}
|
||||
|
||||
ver := windows.RtlGetVersion()
|
||||
if ver == nil {
|
||||
logger.Warnf("vm_os_info metric will miss release info since windows.RtlGetVersion returned nil version")
|
||||
return
|
||||
}
|
||||
os.release = fmt.Sprintf("%d.%d.%d", ver.MajorVersion, ver.MinorVersion, ver.BuildNumber)
|
||||
hostOS.release = fmt.Sprintf("%d.%d.%d", ver.MajorVersion, ver.MinorVersion, ver.BuildNumber)
|
||||
}
|
||||
|
||||
53
lib/appmetrics/unclean_shutdown.go
Normal file
53
lib/appmetrics/unclean_shutdown.go
Normal file
@@ -0,0 +1,53 @@
|
||||
package appmetrics
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/fs"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/logger"
|
||||
)
|
||||
|
||||
// UncleanShutdownMarkerFilename is the marker file used to detect a previous unclean shutdown.
|
||||
const UncleanShutdownMarkerFilename = ".vm_app_running"
|
||||
|
||||
var (
|
||||
uncleanShutdownEnabled atomic.Bool
|
||||
uncleanShutdown uint64
|
||||
)
|
||||
|
||||
// MustCreateUncleanShutdownMarker creates an UncleanShutdownMarkerFilename marker file in the given dirPath.
|
||||
// Must be called once on program startup and paired with a single MustRemoveUncleanShutdownMarker call on exit.
|
||||
//
|
||||
// If the marker file already exists on startup, it indicates a previous unclean shutdown and uncleanShutdown is set to 1.
|
||||
func MustCreateUncleanShutdownMarker(dirPath string) {
|
||||
if !uncleanShutdownEnabled.CompareAndSwap(false, true) {
|
||||
logger.Fatalf("BUG: unclean shutdown marker was already initialized. It could only be called once")
|
||||
}
|
||||
marker := filepath.Join(dirPath, UncleanShutdownMarkerFilename)
|
||||
f, err := os.OpenFile(marker, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600)
|
||||
if err == nil {
|
||||
fs.MustClose(f)
|
||||
return
|
||||
}
|
||||
if os.IsExist(err) {
|
||||
uncleanShutdown = 1
|
||||
logger.Warnf("Previous shutdown was unclean since file %q exists. Please check logs and investigate the reason of unclean shutdown", marker)
|
||||
return
|
||||
}
|
||||
logger.Panicf("FATAL: cannot create unclean shutdown marker %q: %s", marker, err)
|
||||
}
|
||||
|
||||
// MustRemoveUncleanShutdownMarker removes the UncleanShutdownMarkerFilename marker file created by MustCreateUncleanShutdownMarker.
|
||||
// Must be called once, as late as possible before program exit.
|
||||
func MustRemoveUncleanShutdownMarker(dirPath string) {
|
||||
if !uncleanShutdownEnabled.Load() {
|
||||
logger.Fatalf("BUG: unclean shutdown marker was not initialized with MustCreateUncleanShutdownMarker call")
|
||||
}
|
||||
marker := filepath.Join(dirPath, UncleanShutdownMarkerFilename)
|
||||
if err := os.Remove(marker); err != nil {
|
||||
logger.Fatalf("FATAL: cannot remove unclean shutdown marker %q: %s", marker, err)
|
||||
}
|
||||
fs.MustSyncPath(dirPath)
|
||||
}
|
||||
60
lib/appmetrics/unclean_shutdown_test.go
Normal file
60
lib/appmetrics/unclean_shutdown_test.go
Normal file
@@ -0,0 +1,60 @@
|
||||
package appmetrics
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestUncleanShutdownLifecycle(t *testing.T) {
|
||||
|
||||
t.Cleanup(func() {
|
||||
uncleanShutdownEnabled.Store(false)
|
||||
})
|
||||
dirPath := t.TempDir()
|
||||
markerPath := filepath.Join(dirPath, UncleanShutdownMarkerFilename)
|
||||
|
||||
// unclean logic is disabled. the unclean shutdown metric should not be exposed
|
||||
var bb bytes.Buffer
|
||||
writePrometheusMetrics(&bb)
|
||||
if strings.Contains(bb.String(), "vm_app_prev_shutdown_unclean") {
|
||||
t.Fatalf("unexpected unclean shutdown metric before starting the marker")
|
||||
}
|
||||
|
||||
// clean start, the metric must report 0
|
||||
MustCreateUncleanShutdownMarker(dirPath)
|
||||
mustContainUncleanShutdownMetric(t, 0)
|
||||
if _, err := os.Stat(markerPath); err != nil {
|
||||
t.Fatalf("cannot stat the running marker after the first start: %s", err)
|
||||
}
|
||||
MustRemoveUncleanShutdownMarker(dirPath)
|
||||
if _, err := os.Stat(markerPath); !os.IsNotExist(err) {
|
||||
t.Fatalf("unexpected running marker after a clean shutdown; got error %v; want os.ErrNotExist", err)
|
||||
}
|
||||
uncleanShutdownEnabled.Store(false)
|
||||
|
||||
// simulate prev unclean shutdown, the metric must report 1
|
||||
if err := os.WriteFile(markerPath, nil, 0600); err != nil {
|
||||
t.Fatalf("cannot create test marker: %s", err)
|
||||
}
|
||||
MustCreateUncleanShutdownMarker(dirPath)
|
||||
mustContainUncleanShutdownMetric(t, 1)
|
||||
MustRemoveUncleanShutdownMarker(dirPath)
|
||||
if _, err := os.Stat(markerPath); !os.IsNotExist(err) {
|
||||
t.Fatalf("unexpected running marker after a clean shutdown; got error %v; want os.ErrNotExist", err)
|
||||
}
|
||||
}
|
||||
|
||||
func mustContainUncleanShutdownMetric(t *testing.T, value uint64) {
|
||||
t.Helper()
|
||||
|
||||
var bb bytes.Buffer
|
||||
writePrometheusMetrics(&bb)
|
||||
want := "vm_app_prev_shutdown_unclean " + strconv.FormatUint(value, 10) + "\n"
|
||||
if !strings.Contains(bb.String(), want) {
|
||||
t.Fatalf("missing %q in the exported app metrics", want)
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/appmetrics"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/backup/backupnames"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/logger"
|
||||
)
|
||||
@@ -107,7 +108,7 @@ func appendFilesInternal(dst []string, d *os.File) ([]string, error) {
|
||||
}
|
||||
|
||||
func isSpecialFile(name string) bool {
|
||||
return name == "flock.lock" || name == backupnames.RestoreInProgressFilename || name == backupnames.RestoreMarkFileName || strings.HasSuffix(name, ".tmp")
|
||||
return name == "flock.lock" || name == appmetrics.UncleanShutdownMarkerFilename || name == backupnames.RestoreInProgressFilename || name == backupnames.RestoreMarkFileName || strings.HasSuffix(name, ".tmp")
|
||||
}
|
||||
|
||||
// RemoveEmptyDirs recursively removes empty directories under the given dir.
|
||||
|
||||
@@ -39,8 +39,29 @@ func NewArrayBool(name, description string) *ArrayBool {
|
||||
}
|
||||
|
||||
// NewArrayInt returns new ArrayInt with the given name, defaultValue and description.
|
||||
//
|
||||
// -help shows defaultValue as a plain number. Use NewArrayIntWithDynamicDefault when
|
||||
// defaultValue is calculated at runtime.
|
||||
func NewArrayInt(name string, defaultValue int, description string) *ArrayInt {
|
||||
description += fmt.Sprintf(" (default %d)", defaultValue)
|
||||
return newArrayInt(name, defaultValue, strconv.Itoa(defaultValue), description)
|
||||
}
|
||||
|
||||
// NewArrayIntWithDynamicDefault returns new ArrayInt with the given name, defaultValue and description.
|
||||
//
|
||||
// Use it instead of NewArrayInt when defaultValue is calculated at runtime.
|
||||
// See NewIntWithDynamicDefault for why such a value needs a hint.
|
||||
func NewArrayIntWithDynamicDefault(name string, defaultValue int, defaultValueHint, description string) *ArrayInt {
|
||||
if defaultValueHint == "" {
|
||||
panic(fmt.Sprintf("BUG: missing defaultValueHint for -%s", name))
|
||||
}
|
||||
return newArrayInt(name, defaultValue, fmt.Sprintf("%d = %s", defaultValue, defaultValueHint), description)
|
||||
}
|
||||
|
||||
// newArrayInt registers an int array flag, which shows defaultValueText as its default in -help.
|
||||
//
|
||||
// Array flags keep the default in the description, since flag.Var hides an empty DefValue.
|
||||
func newArrayInt(name string, defaultValue int, defaultValueText, description string) *ArrayInt {
|
||||
description += fmt.Sprintf(" (default %s)", defaultValueText)
|
||||
description += "\nSupports `array` of values separated by comma or specified via multiple flags."
|
||||
description += "\nEmpty values are set to default value."
|
||||
a := &ArrayInt{
|
||||
|
||||
@@ -7,6 +7,25 @@ import (
|
||||
"strings"
|
||||
)
|
||||
|
||||
// NewIntWithDynamicDefault returns a new int flag with the given name, defaultValue and description.
|
||||
//
|
||||
// Use it instead of flag.Int when defaultValue is calculated at runtime, for example
|
||||
// from the number of CPU cores. Such a value differs per machine, so -help shows both
|
||||
// the value and defaultValueHint, for example "16 = 2 * availableCPUs".
|
||||
//
|
||||
// Only -help output changes. The flag value stays defaultValue.
|
||||
func NewIntWithDynamicDefault(name string, defaultValue int, defaultValueHint, description string) *int {
|
||||
if defaultValueHint == "" {
|
||||
panic(fmt.Sprintf("BUG: missing defaultValueHint for -%s", name))
|
||||
}
|
||||
p := flag.Int(name, defaultValue, description)
|
||||
|
||||
// DefValue is only the text shown by -help: "default value (as text); for usage message".
|
||||
flag.Lookup(name).DefValue = fmt.Sprintf("%d = %s", defaultValue, defaultValueHint)
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
// WriteFlags writes all the explicitly set flags to w.
|
||||
func WriteFlags(w io.Writer) {
|
||||
flag.Visit(func(f *flag.Flag) {
|
||||
|
||||
51
lib/flagutil/flag_test.go
Normal file
51
lib/flagutil/flag_test.go
Normal file
@@ -0,0 +1,51 @@
|
||||
package flagutil
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// The flags are registered at package level, since flag registration panics when it repeats.
|
||||
var (
|
||||
fooFlagIntDynamicDefault = NewIntWithDynamicDefault("fooFlagIntDynamicDefault", 42, "2 * availableCPUs", "test")
|
||||
fooFlagArrayIntDynamicDefault = NewArrayIntWithDynamicDefault("fooFlagArrayIntDynamicDefault", 42, "2 * availableCPUs", "test")
|
||||
fooFlagArrayIntPlainDefault = NewArrayInt("fooFlagArrayIntPlainDefault", 42, "test")
|
||||
)
|
||||
|
||||
func TestNewIntWithDynamicDefaultSuccess(t *testing.T) {
|
||||
// -help must show the value together with the hint.
|
||||
f := flag.Lookup("fooFlagIntDynamicDefault")
|
||||
if f.DefValue != "42 = 2 * availableCPUs" {
|
||||
t.Fatalf("unexpected DefValue; got %q; want %q", f.DefValue, "42 = 2 * availableCPUs")
|
||||
}
|
||||
|
||||
// the flag value must stay the calculated one.
|
||||
if *fooFlagIntDynamicDefault != 42 {
|
||||
t.Fatalf("unexpected flag value; got %d; want %d", *fooFlagIntDynamicDefault, 42)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewArrayIntWithDynamicDefaultSuccess(t *testing.T) {
|
||||
// array flags keep the default in the description, so the hint must go there.
|
||||
f := flag.Lookup("fooFlagArrayIntDynamicDefault")
|
||||
if !strings.Contains(f.Usage, "(default 42 = 2 * availableCPUs)") {
|
||||
t.Fatalf("missing the hint in the flag description; got %q", f.Usage)
|
||||
}
|
||||
|
||||
// the default value must stay the calculated one.
|
||||
if n := fooFlagArrayIntDynamicDefault.GetOptionalArg(0); n != 42 {
|
||||
t.Fatalf("unexpected default value; got %d; want %d", n, 42)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewArrayIntKeepsPlainDefault(t *testing.T) {
|
||||
// NewArrayInt must keep showing a plain number, since it shares the body with the dynamic one.
|
||||
f := flag.Lookup("fooFlagArrayIntPlainDefault")
|
||||
if !strings.Contains(f.Usage, "(default 42)") {
|
||||
t.Fatalf("unexpected flag description; got %q", f.Usage)
|
||||
}
|
||||
if n := fooFlagArrayIntPlainDefault.GetOptionalArg(0); n != 42 {
|
||||
t.Fatalf("unexpected default value; got %d; want %d", n, 42)
|
||||
}
|
||||
}
|
||||
@@ -1,14 +1,15 @@
|
||||
package fsutil
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"sync"
|
||||
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/cgroup"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/flagutil"
|
||||
)
|
||||
|
||||
var maxConcurrency = flag.Int("fs.maxConcurrency", getDefaultConcurrency(), "The maximum number of concurrent goroutines to work with files; smaller values may help reducing Go scheduling latency "+
|
||||
"on systems with small number of CPU cores; higher values may help reducing data ingestion latency on systems with high-latency storage such as NFS or Ceph")
|
||||
var maxConcurrency = flagutil.NewIntWithDynamicDefault("fs.maxConcurrency", getDefaultConcurrency(), "fsutil.getDefaultConcurrency()",
|
||||
"The maximum number of concurrent goroutines to work with files; smaller values may help reducing Go scheduling latency "+
|
||||
"on systems with small number of CPU cores; higher values may help reducing data ingestion latency on systems with high-latency storage such as NFS or Ceph")
|
||||
|
||||
func getDefaultConcurrency() int {
|
||||
n := min(16*cgroup.AvailableCPUs(), 256)
|
||||
|
||||
@@ -10,16 +10,18 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/cgroup"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/flagutil"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/httpserver"
|
||||
"github.com/VictoriaMetrics/VictoriaMetrics/lib/timerpool"
|
||||
"github.com/VictoriaMetrics/metrics"
|
||||
)
|
||||
|
||||
var (
|
||||
maxConcurrentInserts = flag.Int("maxConcurrentInserts", 2*cgroup.AvailableCPUs(), "The maximum number of concurrent insert requests. "+
|
||||
"Set higher value when clients send data over slow networks. "+
|
||||
"Default value depends on the number of available CPU cores. It should work fine in most cases since it minimizes resource usage. "+
|
||||
"See also -insert.maxQueueDuration")
|
||||
maxConcurrentInserts = flagutil.NewIntWithDynamicDefault("maxConcurrentInserts", 2*cgroup.AvailableCPUs(), "2*cgroup.AvailableCPUs()",
|
||||
"The maximum number of concurrent insert requests. "+
|
||||
"Set higher value when clients send data over slow networks. "+
|
||||
"Default value depends on the number of available CPU cores. It should work fine in most cases since it minimizes resource usage. "+
|
||||
"See also -insert.maxQueueDuration")
|
||||
maxQueueDuration = flag.Duration("insert.maxQueueDuration", time.Minute, "The maximum duration to wait in the queue when -maxConcurrentInserts "+
|
||||
"concurrent insert requests are executed")
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user