InfluxDB 2.7.0异步写入时如何检测数据库状态并抛出异常
解决InfluxDB 2.7.0异步写入时数据库未启动不抛异常的问题
你遇到的问题是因为异步写入API(ASYNCHRONOUS)会把写入任务放到后台队列执行,不会立即阻塞等待请求结果,所以数据库未启动时的连接错误不会在调用write()时直接抛出,而是在后台处理时才会触发,但默认没有暴露这些错误。下面提供两种可行的解决方法:
方法1:写入前主动检测数据库状态
在执行写入操作前,调用客户端的health()方法同步检查数据库健康状态,这个方法会直接向数据库发起请求,一旦数据库未启动或者状态异常,会立刻抛出连接异常。
修改后的代码:
from influxdb_client.client.warnings import MissingPivotFunction from influxdb_client import InfluxDBClient from influxdb_client.client.write_api import ASYNCHRONOUS import warnings warnings.simplefilter("ignore", MissingPivotFunction) client = InfluxDBClient(url="http://localhost:8086", token="my_token", org="test") # 主动检查数据库健康状态 try: health_info = client.health() if health_info.status != "pass": raise RuntimeError(f"数据库状态异常: {health_info.message}") except Exception as e: print(f"数据库连接失败: {str(e)}") raise # 抛出异常终止程序 write_api = client.write_api(write_options=ASYNCHRONOUS) write_api.write("my_bucket", "test", [("test,location=test0 test=1").encode()]) # 异步写入需先等待任务完成再关闭客户端,避免请求丢失 write_api.flush() write_api.close() client.close() print("done")
方法2:给异步写入配置错误回调
通过write_api的on_error参数设置错误回调函数,当后台写入出现异常时,回调函数会捕获到错误,你可以在回调里抛出异常或者做其他处理。
修改后的代码:
from influxdb_client.client.warnings import MissingPivotFunction from influxdb_client import InfluxDBClient from influxdb_client.client.write_api import ASYNCHRONOUS import warnings warnings.simplefilter("ignore", MissingPivotFunction) def handle_write_error(exceptions): # 捕获到写入异常时直接抛出,终止程序 raise RuntimeError(f"写入数据失败: {str(exceptions)}") client = InfluxDBClient(url="http://localhost:8086", token="my_token", org="test") # 初始化write_api时指定错误回调 write_api = client.write_api(write_options=ASYNCHRONOUS, on_error=handle_write_error) write_api.write("my_bucket", "test", [("test,location=test0 test=1").encode()]) # 等待所有异步任务完成 write_api.flush() write_api.close() client.close() print("done")
注意点
- 用异步写入时,一定要先调用
write_api.flush()等待后台任务全部完成,再关闭客户端,否则未处理的请求会被丢弃,错误也无法被捕获。 health()方法需要你的token有足够权限访问数据库的健康检查接口,否则可能返回权限错误,需提前确认token权限。
内容的提问来源于stack exchange,提问作者David01
相关产品推荐
相关产品推荐

