Documentation

使用客户端库对数据进行降采样

查询和下采样存储在InfluxDB中的时间序列数据,并将下采样的数据写回InfluxDB。

本指南使用 PythonInfluxDB 3 Python client library,但您可以使用您选择的运行时和任何可用的 InfluxDB 3 client libraries。本指南还假设您已经 setup your Python project and virtual environment

安装依赖

使用 pip 安装以下依赖项:

  • influxdb_client_3
  • pandas
pip install influxdb3-python pandas

准备 InfluxDB 数据库

降采样过程涉及两个 InfluxDB 数据库。 每个数据库都有一个 保留期限 指定数据在数据库中保留的时间,直到过期并被删除。 通过使用两个数据库,您可以在一个保留期限较短的数据库中存储未修改的高分辨率数据, 然后在一个保留期限较长的数据库中存储降采样的低分辨率数据。

确保您为以下每项准备了一个数据库:

  • 从中查询未修改的数据
  • 另一个用于写入下采样数据的

有关创建数据库的信息,请参见 创建数据库

创建 InfluxDB 客户端

influxdb_client_3模块中使用InfluxDBClient3函数来实例化两个InfluxDB客户端:

  • 一个配置为连接到您的 InfluxDB 数据库并具有 未修改 数据。
  • 另一个配置用于连接到您想要写入的InfluxDB数据库 降采样 数据。

为每个客户提供以下凭证:

  • host: InfluxDB Cloud 专用集群 URL (不包括协议)
  • token: InfluxDB数据库令牌 具有对您想要查询和写入的数据库的读取和写入权限。
  • database: InfluxDB 数据库名称
from influxdb_client_3 import InfluxDBClient3
import pandas

# Instantiate an InfluxDBClient3 client configured for your unmodified database
influxdb_raw = InfluxDBClient3(
    host='cluster-id.a.influxdb.io',
    token='
DATABASE_TOKEN
'
,
database='
RAW_DATABASE_NAME
'
) # Instantiate an InfluxDBClient3 client configured for your downsampled database. # When writing, the org= argument is required by the client (but ignored by InfluxDB). influxdb_downsampled = InfluxDBClient3( host='cluster-id.a.influxdb.io', token='
DATABASE_TOKEN
'
,
database='
DOWNSAMPLED_DATABASE_NAME
'
,
org='' )

查询 InfluxDB

定义一个执行基于时间的聚合的查询

用于下采样时间序列数据最常见的方法是在时间间隔上执行聚合或选择器操作。例如,返回查询时间范围内每小时的平均值。

使用 SQL 或 InfluxQL 通过对时间间隔应用聚合或选择函数来降采样数据。

  1. SELECT 子句中:

  2. 在你的 SELECT 子句中包含一个 GROUP BY 子句,该子句基于来自 DATE_BIN 函数返回的时间间隔进行分组,以及任何其他被查询的标签。下面的示例使用 GROUP BY 1SELECT 子句中的第一列进行分组。

  3. 包含一个 ORDER BY 子句,用于按 time 排序数据。

有关更多信息,请参见 使用SQL聚合数据 - 通过应用基于区间的聚合来调整数据采样.

SELECT
  DATE_BIN(INTERVAL '1 hour', time) AS time,
  room,
  AVG(temp) AS temp,
  AVG(hum) AS hum,
  AVG(co) AS co
FROM home
--In WHERE, time refers to <source_table>.time
WHERE time >= now() - INTERVAL '24 hours'
--1 refers to the DATE_BIN column
GROUP BY 1, room
ORDER BY time
  1. SELECT子句中,对查询字段应用聚合选择器函数。

  2. 包含一个 GROUP BY 子句,该子句按指定的时间间隔对 time() 进行分组。

SELECT
  MEAN(temp) AS temp,
  MEAN(hum) AS hum,
  MEAN(co) AS co
FROM home
WHERE time >= now() - 24h
GROUP BY time(1h)

执行查询

  1. 将查询字符串分配给一个变量。

  2. 使用您实例化的客户端query方法从InfluxDB查询原始数据。提供以下参数。

    • query: 要执行的查询字符串
    • 语言: sqlinfluxql
  3. 使用 to_pandas 方法将返回的 Arrow 表转化为 Pandas 数据框。

# ...

query = '''
SELECT
  DATE_BIN(INTERVAL '1 hour', time) AS time,
  room,
  AVG(temp) AS temp,
  AVG(hum) AS hum,
  AVG(co) AS co
FROM home
--In WHERE, time refers to <source_table>.time
WHERE time >= now() - INTERVAL '24 hours'
--1 refers to the DATE_BIN column
GROUP BY 1, room
ORDER BY 1
'''

table = influxdb_raw.query(query=query, language="sql")
data_frame = table.to_pandas()
# ...

query = '''
SELECT
  MEAN(temp) AS temp,
  MEAN(hum) AS hum,
  MEAN(co) AS co
FROM home
WHERE time >= now() - 24h
GROUP BY time(1h)
'''

table = influxdb_raw.query(query=query, language="influxql")
data_frame = table.to_pandas()
\

将降采样的数据写回到InfluxDB

  1. 对于 InfluxQL 查询结果,在将数据写回 InfluxDB 之前删除 (drop) iox::measurement。您可以避免在稍后查询降采样数据时发生测量名称冲突。

  2. 使用 sort_values 方法按 time 对 Pandas DataFrame 中的数据进行排序,以确保回写到 InfluxDB 的性能尽可能高。

  3. 使用您的 write 方法实例化的下采样客户端将查询结果写回到您的 InfluxDB 数据库中,以获取下采样数据。包括以下参数:

    • record: 包含降采样数据的 Pandas DataFrame
    • data_frame_measurement_name: 目标测量名称
    • data_frame_timestamp_column: 包含每个时间点的时间戳的列
    • data_frame_tag_columns: 标签的 tag 列表

    未在 data_frame_tag_columnsdata_frame_timestamp_column 参数中列出的列将作为 字段 写入 InfluxDB。

# ...

data_frame = data_frame.sort_values(by="time")

influxdb_downsampled.write(
    record=data_frame,
    data_frame_measurement_name="home_ds",
    data_frame_timestamp_column="time",
    data_frame_tag_columns=['room']
)

完整的下采样脚本

from influxdb_client_3 import InfluxDBClient3
import pandas

influxdb_raw = InfluxDBClient3(
    host='cluster-id.a.influxdb.io',
    token='
DATABASE_TOKEN
'
,
database='
RAW_DATABASE_NAME
'
) # When writing, the org= argument is required by the client (but ignored by InfluxDB). influxdb_downsampled = InfluxDBClient3( host='cluster-id.a.influxdb.io', token='
DATABASE_TOKEN
'
,
database='
DOWNSAMPLED_DATABASE_NAME
'
,
org='' ) query = ''' SELECT DATE_BIN(INTERVAL '1 hour', time) AS time, room, AVG(temp) AS temp, AVG(hum) AS hum, AVG(co) AS co FROM home --In WHERE, time refers to <source_table>.time WHERE time >= now() - INTERVAL '24 hours' --1 refers to the DATE_BIN column GROUP BY 1, room ORDER BY 1 ''' table = influxdb_raw.query(query=query, language="sql") data_frame = table.to_pandas() data_frame = data_frame.sort_values(by="time") influxdb_downsampled.write( record=data_frame, data_frame_measurement_name="home_ds", data_frame_timestamp_column="time", data_frame_tag_columns=['room'] )
from influxdb_client_3 import InfluxDBClient3
import pandas

influxdb_raw = InfluxDBClient3(
    host='cluster-id.a.influxdb.io',
    token='
DATABASE_TOKEN
'
,
database='
RAW_DATABASE_NAME
'
) # When writing, the org= argument is required by the client (but ignored by InfluxDB). influxdb_downsampled = InfluxDBClient3( host='cluster-id.a.influxdb.io', token='
DATABASE_TOKEN
'
,
database='
DOWNSAMPLED_DATABASE_NAME
'
,
org='' ) query = ''' SELECT MEAN(temp) AS temp, MEAN(hum) AS hum, MEAN(co) AS co FROM home WHERE time >= now() - 24h GROUP BY time(1h) ''' # To prevent naming conflicts when querying downsampled data, # drop the iox::measurement column before writing the data # with the new measurement. data_frame = data_frame.drop(columns=['iox::measurement']) table = influxdb_raw.query(query=query, language="influxql") data_frame = table.to_pandas() data_frame = data_frame.sort_values(by="time") influxdb_downsampled.write( record=data_frame, data_frame_measurement_name="home_ds", data_frame_timestamp_column="time", data_frame_tag_columns=['room'] )


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

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