使用客户端库对数据进行降采样
查询和下采样存储在InfluxDB中的时间序列数据,并将下采样的数据写回InfluxDB。
本指南使用 Python 和 InfluxDB 3 Python client library,但您可以使用您选择的运行时和任何可用的 InfluxDB 3 client libraries。本指南还假设您已经 setup your Python project and virtual environment。
安装依赖
使用 pip 安装以下依赖项:
influxdb_client_3pandas
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 通过对时间间隔应用聚合或选择函数来降采样数据。
在
SELECT子句中:- 使用
DATE_BIN根据行的时间戳将每一行分配到一个区间,并更新time列为分配的区间时间戳。 您还可以使用DATE_BIN_GAPFILL填补因没有数据而产生的区间间隙 (参见 使用SQL填补数据中的空缺)。 - 对每个查询字段应用一个 aggregate 或 selector 函数。
- 使用
在你的
SELECT子句中包含一个GROUP BY子句,该子句基于来自DATE_BIN函数返回的时间间隔进行分组,以及任何其他被查询的标签。下面的示例使用GROUP BY 1按SELECT子句中的第一列进行分组。包含一个
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
SELECT
MEAN(temp) AS temp,
MEAN(hum) AS hum,
MEAN(co) AS co
FROM home
WHERE time >= now() - 24h
GROUP BY time(1h)
执行查询
将查询字符串分配给一个变量。
使用您实例化的客户端的
query方法从InfluxDB查询原始数据。提供以下参数。- query: 要执行的查询字符串
- 语言:
sql或influxql
使用
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
对于 InfluxQL 查询结果,在将数据写回 InfluxDB 之前删除 (
drop)iox::measurement列 。您可以避免在稍后查询降采样数据时发生测量名称冲突。使用
sort_values方法按time对 Pandas DataFrame 中的数据进行排序,以确保回写到 InfluxDB 的性能尽可能高。使用您的
write方法实例化的下采样客户端将查询结果写回到您的 InfluxDB 数据库中,以获取下采样数据。包括以下参数:- record: 包含降采样数据的 Pandas DataFrame
- data_frame_measurement_name: 目标测量名称
- data_frame_timestamp_column: 包含每个时间点的时间戳的列
- data_frame_tag_columns: 标签的 tag 列表
未在 data_frame_tag_columns 或 data_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']
)