fix(monitor): support influxql derivative and difference function translation (#24749)

This commit is contained in:
Zexi Li
2026-04-30 13:11:58 +08:00
committed by GitHub
parent b46f53a033
commit 71dd2068c4
4 changed files with 61 additions and 6 deletions

2
go.mod
View File

@@ -83,7 +83,7 @@ require (
github.com/vmihailenco/msgpack v4.0.4+incompatible
github.com/xuri/excelize/v2 v2.7.1
github.com/zeebo/xxh3 v1.0.2
github.com/zexi/influxql-to-metricsql v0.1.2
github.com/zexi/influxql-to-metricsql v0.1.3
go.etcd.io/etcd/api/v3 v3.5.7
go.etcd.io/etcd/client/v3 v3.5.7
golang.org/x/crypto v0.41.0

4
go.sum
View File

@@ -1203,8 +1203,8 @@ github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ=
github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0=
github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0=
github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA=
github.com/zexi/influxql-to-metricsql v0.1.2 h1:Akkn1PLuNWjwt+Aqxf1fV+Gg+KUIgouVwtzm4XWy/5w=
github.com/zexi/influxql-to-metricsql v0.1.2/go.mod h1:JlC5FY+6De9ZPxG47G5DOgva8P9X1VaKS4ExzCmhSCc=
github.com/zexi/influxql-to-metricsql v0.1.3 h1:sHYg7bXEqFinV67OQVb5BtzFCzmz9+yRTWfyJs2pOMs=
github.com/zexi/influxql-to-metricsql v0.1.3/go.mod h1:JlC5FY+6De9ZPxG47G5DOgva8P9X1VaKS4ExzCmhSCc=
github.com/zexi/promql/v2 v2.12.1 h1:crHKpULdLLsBZ9b78Rg6qQkugzlk6BHeCj93tw/F5RU=
github.com/zexi/promql/v2 v2.12.1/go.mod h1:2UtzWZGmth95n2qdIZWpJ5yQ0cE5hEvz0dGsajI1Sqg=
go.etcd.io/bbolt v1.3.7 h1:j+zJOnnEjF/kyHlDDgGnVL/AIqIJPq8UoB2GSNfkUfQ=

View File

@@ -415,6 +415,44 @@ func getAggrExpr(ops []*AggrOperator, expr promql.Expr) promql.Expr {
promql.Expressions{
aggrOp.Args[0],
restExpr})
case "non_negative_derivative":
// InfluxQL: non_negative_derivative(mean("field"), 1s) computes per-second non-negative rate of change
// MetricsQL: rate() is the equivalent for counter-like metrics
// Unwrap inner aggregation (e.g., avg_over_time) and apply rate() to the base metric
if callExpr, ok := restExpr.(*promql.Call); ok && len(callExpr.Args) > 0 {
expr = newAggrExpr("rate", promql.ValueTypeMatrix, promql.ValueTypeVector, callExpr.Args[0])
} else {
expr = newAggrExpr("rate", promql.ValueTypeMatrix, promql.ValueTypeVector, restExpr)
}
case "derivative":
// InfluxQL: derivative(mean("field"), 1s) computes per-second rate of change (can be negative)
// MetricsQL: deriv() is the closest equivalent
if callExpr, ok := restExpr.(*promql.Call); ok && len(callExpr.Args) > 0 {
expr = newAggrExpr("deriv", promql.ValueTypeMatrix, promql.ValueTypeVector, callExpr.Args[0])
} else {
expr = newAggrExpr("deriv", promql.ValueTypeMatrix, promql.ValueTypeVector, restExpr)
}
case "difference":
// InfluxQL: difference(mean("field")) computes difference between consecutive points
// MetricsQL: delta() is the closest equivalent
if callExpr, ok := restExpr.(*promql.Call); ok && len(callExpr.Args) > 0 {
expr = newAggrExpr("delta", promql.ValueTypeMatrix, promql.ValueTypeVector, callExpr.Args[0])
} else {
expr = newAggrExpr("delta", promql.ValueTypeMatrix, promql.ValueTypeVector, restExpr)
}
case "non_negative_difference":
// InfluxQL: non_negative_difference() - like difference but only non-negative values
// MetricsQL: increase() is the closest equivalent
if callExpr, ok := restExpr.(*promql.Call); ok && len(callExpr.Args) > 0 {
expr = newAggrExpr("increase", promql.ValueTypeMatrix, promql.ValueTypeVector, callExpr.Args[0])
} else {
expr = newAggrExpr("increase", promql.ValueTypeMatrix, promql.ValueTypeVector, restExpr)
}
case "elapsed":
// elapsed is not directly supported, pass through
case "moving_average":
// moving_average is not directly supported, pass through as avg_over_time
expr = newAggrExpr("avg_over_time", promql.ValueTypeMatrix, promql.ValueTypeVector, restExpr)
}
return expr
}
@@ -431,7 +469,7 @@ func newAggrOperatorByName(name string) *AggrOperator {
}
func getAggrOperator(op *influxql.Call) ([]*AggrOperator, error) {
if len(op.Args) != 1 && !MUL_ARGS_AGGREGATOR.Has(op.Name) {
if len(op.Args) != 1 && !MUL_ARGS_AGGREGATOR.Has(op.Name) && !hasDurationLiteralExtraArgs(op) {
return nil, errors.Errorf("not supported aggregator: %s with args: %#v", op.String(), op.Args)
}
aggOp := newAggrOperatorByName(op.Name)
@@ -548,8 +586,25 @@ func trimRegexDelimiters(metricName string) string {
return prefix + regexPart
}
// hasDurationLiteralExtraArgs checks if a call has extra args that are all DurationLiterals or IntegerLiterals.
// This handles functions like non_negative_derivative(mean("field"), 1s) where 1s is a duration parameter.
func hasDurationLiteralExtraArgs(c *influxql.Call) bool {
if len(c.Args) <= 1 {
return false
}
for _, arg := range c.Args[1:] {
switch arg.(type) {
case *influxql.DurationLiteral, *influxql.IntegerLiteral:
continue
default:
return false
}
}
return true
}
func getCallVariable(c *influxql.Call) (string, error) {
if len(c.Args) != 1 && !MUL_ARGS_AGGREGATOR.Has(c.Name) {
if len(c.Args) != 1 && !MUL_ARGS_AGGREGATOR.Has(c.Name) && !hasDurationLiteralExtraArgs(c) {
return "", errors.Errorf("length of call %q args %#v != 1", c.Name, c.Args)
}
switch args := c.Args[0].(type) {

2
vendor/modules.txt vendored
View File

@@ -1846,7 +1846,7 @@ github.com/yusufpapurcu/wmi
# github.com/zeebo/xxh3 v1.0.2
## explicit; go 1.17
github.com/zeebo/xxh3
# github.com/zexi/influxql-to-metricsql v0.1.2
# github.com/zexi/influxql-to-metricsql v0.1.3
## explicit; go 1.18
github.com/zexi/influxql-to-metricsql/converter
github.com/zexi/influxql-to-metricsql/converter/translator