Documentation

来自节点

from节点选择通过StreamNode.流动的数据的一个子集。流节点允许您选择要处理的流的哪个部分。

示例:

stream
  |from()
    .database('mydb')
    .retentionPolicy('myrp')
    .measurement('mymeasurement')
    .where(lambda: "host" =~ /logger\d+/)
  |window()
  ...

上述示例仅从数据库 mydb 和保留策略 myrp 以及测量 mymeasurement 中选择数据点,条件是标签 host 匹配正则表达式 logger\d+

构造函数

链式方法描述
from ( )创建一个新的流节点,可以通过数据库、保留策略、测量和条件属性进一步过滤。可以多次调用 From 来创建数据流的多个独立分支。

属性方法

设置器描述
database ( value string)数据库名称。如果为空,将使用任何数据库。
groupBy ( tag ...interface{})通过一组标签对数据进行分组。
groupByMeasurement ( )如果设置,将在组 ID 中包含测量名称。以及任何其他的分组维度。
measurement ( value string)测量名称,如果为空,将使用任何测量。
quiet ( )抑制此节点的所有错误日志事件。
retentionPolicy ( value string)保留策略名称,如果为空,则使用任何保留策略。
round ( value time.Duration)可选的持续时间,用于对时间戳进行四舍五入。 有助于确保数据点落在特定的边界上示例:流
truncate ( value time.Duration)截断时间戳的可选持续时间。帮助确保数据点落在特定边界上 示例:stream
where ( lambda ast.LambdaNode)使用给定的表达式过滤当前流。该表达式是Kapacitor表达式。Kapacitor表达式是InfluxQL WHERE表达式的超集。有关更多信息,请参阅expression文档。

链式调用方法

警报, 障碍, 底部, 变更检测, 组合, 计数, 累积和, 死人开关, 默认, 删除, 导数, 差异, 唯一, Ec2自动缩放, 经过时间, 评估, 第一个, 展平, 来自, 霍尔特-温特斯, 霍尔特-温特斯带拟合, Http输出, Http Post, InfluxDB输出, 连接, K8s自动缩放, Kapacitor回环, 最后, 日志, 最大值, 均值, 中位数, 最小值, 众数, 移动平均, 百分位数, 样本, 移位, 侧载, 扩展, 状态计数, 状态持续时间, 统计, 标准差, 总和, 群体自动缩放, 顶部, 联合, 窗口


属性

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

数据库

数据库名称。
如果为空,将使用任意数据库。

from.database(value string)

分组

按一组标签对数据进行分组。

可以传递字面量 * 来按所有维度分组。

示例:

  stream
      |from()
          .groupBy(*)
from.groupBy(tag ...interface{})

按测量分组

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

示例:

 stream
      |from()
          .database('mydb')
          .groupByMeasurement()
          .groupBy('host')

上述示例从数据库‘mydb’中选择所有测量值,然后每个点按主机标签和测量名称进行分组。因此将测量值保留在各自的组中。

from.groupByMeasurement()

测量

测量名称
如果为空,将使用任何测量。

from.measurement(value string)

安静

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

from.quiet()

保留策略

保留策略名称
如果为空,将使用任何保留策略。

from.retentionPolicy(value string)

四舍五入

可选的时长用于四舍五入时间戳。 有助于确保数据点落在特定的边界上 示例:

    stream
       |from()
           .measurement('mydata')
           .round(1s)

所有传入的数据将四舍五入到最接近的1秒边界。

from.round(value time.Duration)

截断

可选的持续时间用于截断时间戳。
有助于确保数据点落在特定边界上。
示例:

    stream
       |from()
           .measurement('mydata')
           .truncate(1s)

所有传入的数据将被截断到1秒的分辨率。

from.truncate(value time.Duration)

在哪里

使用给定的表达式过滤当前流。 该表达式是一个Kapacitor表达式。Kapacitor 表达式是InfluxQL WHERE表达式的超集。 有关更多信息,请参见expression文档。

多次调用 Where 方法将把每个表达式用 AND 连接在一起。

示例:

    stream
       |from()
          .where(lambda: condition1)
          .where(lambda: condition2)

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

    stream
       |from()
          .where(lambda: condition1 AND condition2)

注意:如果您想要多个不同的数据流,请务必始终使用 |from

示例:

  var data = stream
      |from()
          .measurement('cpu')
  var total = data
      .where(lambda: "cpu" == 'cpu-total')
  var others = data
      .where(lambda: "cpu" != 'cpu-total')

上面的示例等价于下面的示例,这显然不是预期的结果。

示例:

  var data = stream
      |from()
          .measurement('cpu')
          .where(lambda: "cpu" == 'cpu-total' AND "cpu" != 'cpu-total')
  var total = data
  var others = total

下面的示例将创建两个不同的流,每个流选择原始流的不同子集。

示例:

  var data = stream
      |from()
          .measurement('cpu')
  var total = stream
      |from()
          .measurement('cpu')
          .where(lambda: "cpu" == 'cpu-total')
  var others = stream
      |from()
          .measurement('cpu')
          .where(lambda: "cpu" != 'cpu-total')

如果为空,则所有数据点都被视为匹配。

from.where(lambda ast.LambdaNode)

链式调用方法

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

警告

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

from|alert()

返回: AlertNode

障碍

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

每个时段都会发出一个 BarrierMessage。

from|barrier()

返回: BarrierNode

底部

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

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

返回: InfluxQLNode

变化检测

创建一个新节点,仅在与前一个点不同时发出新点。

from|changeDetect(field string)

返回: ChangeDetectNode

合并

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

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

返回: CombineNode

计数

计算点的数量。

from|count(field string)

返回: InfluxQLNode

累积和

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

from|cumulativeSum(field string)

返回: InfluxQLNode

死者

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

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

示例:

    var data = stream
        |from()...
    // 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 = stream
        |from()...
    // 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 = stream
        |from()...
    // 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 = stream
        |from()...
    // 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...
from|deadman(threshold float64, interval time.Duration, expr ...ast.LambdaNode)

返回: AlertNode

默认

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

from|default()

返回: DefaultNode

删除

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

from|delete()

返回: DeleteNode

导数

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

from|derivative(field string)

返回: DerivativeNode

差异

计算点之间的差异,与经过的时间无关。

from|difference(field string)

返回: InfluxQLNode

不同

生成仅有独特点的批次。

from|distinct(field string)

返回: InfluxQLNode

自动扩展EC2

创建可以触发ec2自动扩展组的自动扩展事件的节点。

from|ec2Autoscale()

返回: Ec2AutoscaleNode

经过时间

计算点之间的经过时间。

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

返回: InfluxQLNode

评估

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

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

返回: EvalNode

第一

选择第一个点。

from|first(field string)

返回: InfluxQLNode

扁平化

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

from|flatten()

返回: FlattenNode

来自

创建一个新的流节点,可以进一步使用数据库、保留策略、测量和条件属性进行过滤。可以多次调用 From 来创建数据流的多个独立分支。

示例:

    // Select the 'cpu' measurement from just the database 'mydb'
    // and retention policy 'myrp'.
    var cpu = stream
        |from()
            .database('mydb')
            .retentionPolicy('myrp')
            .measurement('cpu')
    // Select the 'load' measurement from any database and retention policy.
    var load = stream
        |from()
            .measurement('load')
    // Join cpu and load streams and do further processing.
    cpu
        |join(load)
            .as('cpu', 'load')
        ...
from|from()

返回: FromNode

霍尔特-温特斯

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

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

返回: InfluxQLNode

霍尔特-冬季法与拟合

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

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

返回: InfluxQLNode

Http输出

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

from|httpOut(endpoint string)

返回: HTTPOutNode

HttpPost

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

from|httpPost(url ...string)

返回: HTTPPostNode

InfluxDB输出

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

from|influxDBOut()

返回: InfluxDBOutNode

加入

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

from|join(others ...Node)

返回: JoinNode

K8s自动缩放

创建一个可以触发自动伸缩事件的kubernetes集群节点。

from|k8sAutoscale()

返回: K8sAutoscaleNode

Kapacitor回环

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

from|kapacitorLoopback()

返回: KapacitorLoopbackNode

最后

选择最后一点。

from|last(field string)

返回: InfluxQLNode

日志

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

from|log()

返回: LogNode

最大值

选择最大点。

from|max(field string)

返回: InfluxQLNode

均值

计算数据的均值。

from|mean(field string)

返回: InfluxQLNode

中位数

计算数据的中位数。

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

from|median(field string)

返回: InfluxQLNode

最小值

选择最小点。

from|min(field string)

返回: InfluxQLNode

模式

计算数据的众数。

from|mode(field string)

返回: InfluxQLNode

移动平均

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

from|movingAverage(field string, window int64)

返回: InfluxQLNode

百分位数

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

from|percentile(field string, percentile float64)

返回: InfluxQLNode

示例

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

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

from|sample(rate interface{})

返回: SampleNode

移位

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

from|shift(shift time.Duration)

返回: ShiftNode

侧载

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

from|sideload()

返回: SideloadNode

扩散

计算 minmax 点之间的差。

from|spread(field string)

返回: InfluxQLNode

状态计数

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

from|stateCount(expression ast.LambdaNode)

返回: StateCountNode

状态持续时间

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

from|stateDuration(expression ast.LambdaNode)

返回: StateDurationNode

统计

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

from|stats(interval time.Duration)

返回结果: StatsNode

标准差

计算标准差。

from|stddev(field string)

返回: InfluxQLNode

总和

计算所有值的总和。

from|sum(field string)

返回: InfluxQLNode

群集自动缩放

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

from|swarmAutoscale()

返回: SwarmAutoscaleNode

顶部

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

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

返回: InfluxQLNode

联合

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

from|union(node ...Node)

返回: UnionNode

窗口

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

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

from|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 企业版是建立在核心基础之上的商业版本,增加了历史查询能力、读取副本、高可用性、可扩展性和细粒度安全性。

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