Python InfluxDB2异步write_api.write操作如何判断写入成功?
InfluxDB异步写入结果判断方案
Python官方方案(适配influxdb-client库)
你使用的异步写入模式下,write方法返回的是asyncio.Future类实例,官方提供了标准的结果获取机制,不需要访问受保护属性:
方案1:等待任务完成捕获异常
直接await返回的result对象,写入失败(包含服务宕机、鉴权失败、格式错误等所有异常场景)会直接抛出对应异常,捕获即可判断结果:
from influxdb_client import InfluxDBClient from influxdb_client.client.write_api import ASYNCHRONOUS import asyncio client = InfluxDBClient(url="你的InfluxDB地址", token="你的访问token", org="你的组织名") write_api = client.write_api(write_options=ASYNCHRONOUS) async def write_data(): try: # 等待写入操作执行完成 await write_api.write(bucket=bucket, data_frame_measurement_name=field_key, record=a_data_frame) # 无异常走到此处即为写入成功 print("写入成功") except Exception as e: # 所有写入错误都会在此处捕获 print(f"写入失败,错误信息:{str(e)}") asyncio.run(write_data())
方案2:添加异步回调
如果不需要阻塞等待结果,可以给返回的Future对象绑定成功/失败回调,异步处理结果:
def handle_success(result): print("写入成功") def handle_fail(exc): print(f"写入失败,错误:{str(exc)}") result = write_api.write(bucket=bucket, data_frame_measurement_name=field_key, record=a_data_frame) # 绑定回调函数 result.add_done_callback(lambda f: handle_success(f.result()) if not f.exception() else handle_fail(f.exception()))
批量写入额外配置
如果开启了批量异步写入,可以在初始化时直接配置全局的成功、失败回调,不用单独处理每次写入:
from influxdb_client.client.write_api import WriteOptions write_options = WriteOptions( write_type=ASYNCHRONOUS, success_callback=lambda batch, data: print("批次数据写入成功"), error_callback=lambda batch, exc, data: print(f"批次数据写入失败: {exc}") ) write_api = client.write_api(write_options=write_options)
非Python通用方案
如果使用其他语言客户端,或者需要独立校验写入结果,可以通过查询校验逻辑判断:
- 写入前统计待写入数据的时间范围、总点数
- 写入完成后用Flux语句查询对应bucket、measurement、时间范围内的总点数,和待写入数量对比,一致即可判定写入成功
示例Flux查询语句:
from(bucket: "你的bucket名称") |> range(start: 待写入数据最小时间戳, stop: 待写入数据最大时间戳) |> filter(fn: (r) => r._measurement == "你的measurement名称") |> count()
内容的提问来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

