Documentation

查询节点

query节点定义了用于处理批量数据的源和调度。从InfluxDB查询数据,经过query节点计算后,传递到数据管道中。

示例:

batch
  |query('''
    SELECT mean("value")
    FROM "telegraf"."default".cpu_usage_idle
    WHERE "host" = 'serverA'
  ''')
    .period(1m)
    .every(20s)
    .groupBy(time(10s), 'cpu')
    ...

在上面的例子中,InfluxDB每20秒查询一次;返回的时间窗口跨度为1分钟,并分组为10秒的时间桶。

要在query节点中使用InfluxQL高级语法的函数,您必须在TICKScript中表达部分InfluxQL查询。高级语法计算时间区间内的嵌套聚合,然后对结果应用外部聚合。

例如,以下 InfluxQL 脚本计算每个 cpu 的 10 秒平均 cpu 使用率的非负差异,但对于 query 节点来说是无效的:

SELECT non_negative_difference(mean("value"))
  FROM "telegraf"."default"."cpu_usage_idle"
  WHERE "host" = 'serverA'
  GROUP BY time(10s), "cpu"
  ...

要计算上面的结果对于query节点,你必须使用TICKScript指定分组和外部聚合:

batch
  |query('''
    SELECT mean("value")
    FROM "telegraf"."default".cpu_usage_idle
    WHERE "host" = 'serverA'
  ''')
    .period(1m)
    .every(1m)
    .groupBy(time(10s), 'cpu')
  | difference('max_usage')
  | where(lambda: "difference" >= 0)
  ...

构造函数

链式方法描述
query ( q string)要执行的查询。不得在WHERE子句中包含时间条件或包含GROUP BY子句。时间条件根据周期、偏移量和计划动态添加。GROUP BY子句根据传递给groupBy方法的维度动态添加。

属性方法

设置器描述
align ( )对查询的开始和结束时间进行对齐,使其与QueryNode的每个属性的偶数边界一致。如果使用QueryNode.Cron属性,则不适用。
alignGroup ( )根据查询的开始时间以时间间隔对组进行对齐
cluster ( value string)配置的 InfluxDB 集群的名称。如果为空,将使用默认集群。
cron ( value string)使用cron语法定义计划。
every ( value time.Duration)查询InfluxDB的频率。
fill ( value interface{})填充数据。选项包括:
groupBy ( d ...interface{})通过一组维度对数据进行分组。可以指定一个时间维度。
groupByMeasurement ( )如果设置,将在组 ID 中包含测量名称。以及任何其他的分组维度。
offset ( value time.Duration)从当前时间向后查询的时间跨度
period ( value time.Duration)将从InfluxDB查询的时间段或长度
quiet ( )抑制此节点的所有错误日志事件。

链式调用方法

警报, 障碍, 底部, 变化检测, 组合, 计数, 累积和, 死信阀, 默认, 删除, 导数, 差异, 不同, Ec2自动缩放, 经过时间, 评估, 第一个, 扁平化, 霍尔特-温特斯, 霍尔特-温特斯拟合, Http输出, Http发布, InfluxDB输出, 连接, K8s自动缩放, Kapacitor循环反馈, 最后, 日志, 最大值, 均值, 中位数, 最小值, 众数, 移动平均, 百分位数, 样本, 偏移, 并行加载, 扩散, 状态计数, 状态持续时间, 统计, 标准差, , 群体自动缩放, 顶部, 涓流, 联合, 哪里, 窗口


属性

属性方法修改调用节点的状态。它们不会向管道中添加另一个节点,并始终返回对调用节点的引用。属性方法使用.运算符标记。

对齐

对齐查询的开始和结束时间,使其具有 QueryNode.Every 属性的偶数边界。如果使用 QueryNode.Cron 属性,则不适用。

query.align()

对齐组

将按时间间隔分组与查询的开始时间对齐

query.alignGroup()

聚类

配置的InfluxDB集群的名称。 如果为空,将使用默认集群。

query.cluster(value string)

定时任务

使用cron语法定义一个调度。

具体的cron实现文档在这里: https://github.com/gorhill/cronexpr#implementation

Cron属性与Every属性是互斥的。

query.cron(value string)

每一个

多频繁查询InfluxDB。

每个属性与Cron属性是互斥的。

query.every(value time.Duration)

填充

填写数据。 选项如下:

  • 任何数值
  • null - 表现与默认值相同
  • previous - 报告前一个窗口的值
  • none - 抑制时间戳和值为null的值
  • 线性 - 报告线性插值的结果
query.fill(value interface{})

分组

按一组维度对数据进行分组。可以指定一个时间维度。

此属性向查询添加了一个 GROUP BY 子句,因此在使用 GROUP BY 查询 InfluxDB 时,所有正常行为都适用。

当您的周期长于您的分组时间间隔时,请使用按时间分组。

示例:

    batch
        |query(...)
            .period(1m)
            .every(1m)
            .groupBy(time(10s), 'tag1', 'tag2'))
            .align()

按时间偏移分组也是可能的。

示例:

    batch
        |query(...)
            .period(1m)
            .every(1m)
            .groupBy(time(10s, -5s), 'tag1', 'tag2'))
            .align()
            .offset(5s)

建议将 QueryNode.AlignQueryNode.Offset 与按时间维度分组一起使用,以便时间边界与分组间隔匹配。要自动将分组间隔对齐到查询时间的开始,请使用 QueryNode.AlignGroup. 这在更复杂的情况下很有用,例如当 groupBy 时间段超过查询频率时。

示例:

    batch
        |query(...)
            .period(5m)
            .every(30s)
            .groupBy(time(1m), 'tag1', 'tag2')
            .align()
            .alignGroup()

对于上述示例,如果没有 QueryNode.AlignGroup, Kapacitor 发出的每个其他查询 (在每分钟的 :30)将会对齐到 :00 秒,而不是期望的 :30 秒, 这将创建 6 个分组间隔,而不是 5 个,第一个和最后一个 将只有 30 秒的数据,而不是整整 1 分钟。 如果在与 QueryNode.AlignGroup, 一起使用分组时间偏移(即 time(t, offset)), 则会首先进行对齐,然后会偏移所指定的数量。

注意:由于 QueryNode.Offset 本身是一个负属性,因此“时间”函数的第二个“offset”参数为负值以匹配。

query.groupBy(d ...interface{})

按测量分组

如果设置,将在组 ID 中包含测量名称。与任何其他分组维度一起。

示例:

 batch
      |query('SELECT sum("value") FROM "telegraf"."autogen"./process_.*/')
          .groupByMeasurement()
          .groupBy('host')

上述示例从几个测量中选择数据,这些测量匹配 `/process_.*/`,然后每个点按主机标签和测量名称进行分组。这样可以将测量保持在各自的组中。

query.groupByMeasurement()

偏移量

从当前时间向回查询多远

例如,偏移量为2小时且每次为5分钟,Kapacitor将每5分钟查询一次InfluxDB,获取2小时前的数据窗口。

这也适用于 Cron 调度。如果 cron 指定每星期天凌晨 1 点运行,而时区偏移为 1 小时。那么在星期天凌晨 1 点时,将查询前一天晚上 12 点的数据。

query.offset(value time.Duration)

周期

将从InfluxDB查询的时间段或持续时间

query.period(value time.Duration)

安静

抑制来自此节点的所有错误日志事件。

query.quiet()

链式调用方法

链式方法在调用节点的子节点中创建一个新的节点。它们不会修改调用节点。链式方法使用 | 运算符标记。

警告

创建一个警报节点,可以触发警报。

query|alert()

返回: AlertNode

障碍

创建一个新的障碍节点,它定期发出障碍消息。

每个周期都会发出一条 barrier消息。

query|barrier()

返回: BarrierNode

底部

选择底部 num 点用于 field 并按任何额外标签或字段排序。

query|bottom(num int64, field string, fieldsAndTags ...string)

返回: InfluxQLNode

变更检测

创建一个新节点,只有在与前一个点不同的情况下才发出新点。

query|changeDetect(field string)

返回: ChangeDetectNode

合并

将此节点与自身结合。数据根据时间戳进行结合。

query|combine(expressions ...ast.LambdaNode)

返回: CombineNode

计数

计算点的数量。

query|count(field string)

返回: InfluxQLNode

累积和

计算每个接收到的点的累积和。 每收集到一个点就会发出一个点。

query|cumulativeSum(field string)

返回: InfluxQLNode

死者

用于在低吞吐量时创建警报的辅助函数,也称为死手开关。

  • 阈值:如果吞吐量在点/区间中下降到阈值以下,则触发警报。
  • 间隔:检查吞吐量的频率。
  • 表达式:可选的表达式列表,供评估使用。对于时间警报非常有用。

示例:

    var data = batch
        |query()...
    // Trigger critical alert if the throughput drops below 100 points per 10s and checked every 10s.
    data
        |deadman(100.0, 10s)
    //Do normal processing of data
    data...

上面的内容等同于这个示例:

    var data = batch
        |query()...
    // Trigger critical alert if the throughput drops below 100 points per 10s and checked every 10s.
    data
        |stats(10s)
            .align()
        |derivative('emitted')
            .unit(10s)
            .nonNegative()
        |alert()
            .id('node \'stream0\' in task \'{{ .TaskName }}\'')
            .message('{{ .ID }} is {{ if eq .Level "OK" }}alive{{ else }}dead{{ end }}: {{ index .Fields "emitted" | printf "%0.3f" }} points/10s.')
            .crit(lambda: "emitted" <= 100.0)
    //Do normal processing of data
    data...

可以通过“deadman”配置部分全局配置idmessage警报属性。

由于AlertNode是最后一部分,可以像往常一样进一步修改。 示例:

    var data = batch
        |query()...
    // Trigger critical alert if the throughput drops below 100 points per 10s and checked every 10s.
    data
        |deadman(100.0, 10s)
            .slack()
            .channel('#dead_tasks')
    //Do normal processing of data
    data...

您可以指定额外的lambda表达式,以进一步限制何时触发死手按钮。 示例:

    var data = batch
        |query()...
    // Trigger critical alert if the throughput drops below 100 points per 10s and checked every 10s.
    // Only trigger the alert if the time of day is between 8am-5pm.
    data
        |deadman(100.0, 10s, lambda: hour("time") >= 8 AND hour("time") <= 17)
    //Do normal processing of data
    data...
query|deadman(threshold float64, interval time.Duration, expr ...ast.LambdaNode)

返回: AlertNode

默认

创建一个节点,可以为缺失的标签或字段设置默认值。

query|default()

返回: DefaultNode

删除

创建一个可以删除标签或字段的节点。

query|delete()

返回: DeleteNode

导数

创建一个新节点,该节点计算相邻点的导数。

query|derivative(field string)

返回: DerivativeNode

差异

计算独立于经过时间的点之间的差异。

query|difference(field string)

返回: InfluxQLNode

唯一

生成仅包含不同点的批次。

query|distinct(field string)

返回: InfluxQLNode

Ec2Autoscale

创建一个可以触发自动缩放事件的 EC2 自动缩放组节点。

query|ec2Autoscale()

返回: Ec2AutoscaleNode

经过时间

计算点之间的经过时间。

query|elapsed(field string, unit time.Duration)

返回: InfluxQLNode

评估

创建一个评估节点,该节点将对每个数据点评估给定的变换函数。可以提供表达式列表,并将按给定顺序进行评估。结果可供后续表达式使用。

query|eval(expressions ...ast.LambdaNode)

返回: EvalNode

第一

选择第一个点。

query|first(field string)

返回: InfluxQLNode

扁平化

将具有相似时间的点合并为一个点。

query|flatten()

返回: FlattenNode

霍尔特-温特斯

计算一个数据集的霍尔特-温特斯(/influxdb/v1/query_language/functions/#holt-winters)预测。

query|holtWinters(field string, h int64, m int64, interval time.Duration)

返回: InfluxQLNode

霍尔特-冬季法与拟合

计算Holt-Winters (/influxdb/v1/query_language/functions/#holt-winters) 数据集的预测。 此方法还输出用于拟合数据的所有点,除了预测的数据。

query|holtWintersWithFit(field string, h int64, m int64, interval time.Duration)

返回: InfluxQLNode

Http输出

创建一个HTTP输出节点,用于缓存它所接收到的最新数据。缓存的数据可以在给定的端点访问。该端点是从运行任务的API端点的相对路径。例如,如果任务端点位于 /kapacitor/v1/tasks/<task_id> 且端点为 top10,那么数据可以从 /kapacitor/v1/tasks/<task_id>/top10 请求。

query|httpOut(endpoint string)

返回: HTTPOutNode

HttpPost

创建一个HTTP Post节点,将接收到的数据POST到提供的HTTP端点。HttpPost期望0或1个参数。如果提供0个参数,必须指定一个端点属性方法。

query|httpPost(url ...string)

返回: HTTPPostNode

InfluxDB输出

创建一个 influxdb 输出节点,将传入的数据存储到 InfluxDB 中。

query|influxDBOut()

返回: InfluxDBOutNode

加入

将此节点与其他节点连接。数据是基于时间戳进行连接的。

query|join(others ...Node)

返回: JoinNode

K8s自缩放

创建一个可以触发Kubernetes集群自适应缩放事件的节点。

query|k8sAutoscale()

返回: K8sAutoscaleNode

Kapacitor循环回路

创建一个将数据作为流发送回Kapacitor的kapacitor循环节点。

query|kapacitorLoopback()

返回: KapacitorLoopbackNode

最后

选择最后一点。

query|last(field string)

返回: InfluxQLNode

日志

创建一个节点,记录它接收到的所有数据。

query|log()

返回: LogNode

最大值

选择最大点。

query|max(field string)

返回: InfluxQLNode

均值

计算数据的平均值。

query|mean(field string)

返回: InfluxQLNode

中位数

计算数据的中位数。

注意:此方法不是选择器。如果你想要中位数,请使用 .percentile(field, 50.0)

query|median(field string)

返回: InfluxQLNode

最小值

选择最小点。

query|min(field string)

返回: InfluxQLNode

模式

计算数据的众数。

query|mode(field string)

返回: InfluxQLNode

移动平均

计算最后窗口点的移动平均值。 在窗口填满之前不会发出任何点。

query|movingAverage(field string, window int64)

返回: InfluxQLNode

百分位数

在给定百分位数处选择一个点。 这是一个选择器函数,不执行点之间的插值。

query|percentile(field string, percentile float64)

返回: InfluxQLNode

示例

创建一个新节点,该节点对传入的点或批次进行采样。

每个指定的计数或持续时间将会发出一个点。

query|sample(rate interface{})

返回: SampleNode

移位

创建一个新的节点,按时间移动传入的点或批次。

query|shift(shift time.Duration)

返回: ShiftNode

侧载

创建一个可以从外部源加载数据的节点。

query|sideload()

返回: SideloadNode

扩散

计算 minmax 点之间的差。

query|spread(field string)

返回: InfluxQLNode

状态计数

创建一个节点,用于跟踪给定状态中连续点的数量。

query|stateCount(expression ast.LambdaNode)

返回: StateCountNode

状态持续时间

创建一个跟踪给定状态下持续时间的节点。

query|stateDuration(expression ast.LambdaNode)

返回: StateDurationNode

统计

创建一个新的数据流,其中包含节点的内部统计信息。 间隔表示根据实时多长时间发出一次统计信息。 这意味着间隔时间与源节点接收的数据点次数无关。

query|stats(interval time.Duration)

返回结果: StatsNode

标准差

计算标准差。

query|stddev(field string)

返回: InfluxQLNode

总和

计算所有值的总和。

query|sum(field string)

返回: InfluxQLNode

群集自动缩放

创建一个可以触发Docker swarm集群的自动缩放事件的节点。

query|swarmAutoscale()

返回: SwarmAutoscaleNode

顶部

选择前 num 个点用于 field 并按任何额外标签或字段排序。

query|top(num int64, field string, fieldsAndTags ...string)

返回: InfluxQLNode

涓流

创建一个新的节点,将批量数据转换为流数据。

query|trickle()

返回: TrickleNode

联合

执行该节点与所有其他给定节点的并集。

query|union(node ...Node)

返回: UnionNode

在哪里

创建一个新节点,该节点根据给定的表达式过滤数据流。

query|where(expression ast.LambdaNode)

返回: WhereNode

窗口

创建一个新的节点,通过时间窗口化流。

注意:窗口只能应用于流边缘。

query|window()

返回: WindowNode



Flux的未来

Flux 正在进入维护模式。您可以像现在一样继续使用它,而无需对您的代码进行任何更改。

阅读更多

InfluxDB 3 开源版本现已公开Alpha测试

InfluxDB 3 Open Source is now available for alpha testing, licensed under MIT or Apache 2 licensing.

我们将发布两个产品作为测试版的一部分。

InfluxDB 3 核心,是我们新的开源产品。 它是一个用于时间序列和事件数据的实时数据引擎。 InfluxDB 3 企业版是建立在核心基础之上的商业版本,增加了历史查询能力、读取副本、高可用性、可扩展性和细粒度安全性。

有关如何开始的更多信息,请查看: