You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

InfluxDB批量写入measurement仅随机2条数据落库问题排查

InfluxDB随机仅写入2条数据问题排查

问题复现代码

if __name__ == '__main__':
     measurements = [
       {
            "name": "name_of_my_measurement",
            "tags": {
                "customer_code": "FOO"
            },
            "fields": {
                "data_path": "s3://foo/2022-05-22",
                "count": "40690526"
            }
        },
        {
            "name": "name_of_my_measurement",
            "tags": {
                "customer_code": "FOO"
            },
            "fields": {
                "data_path": "s3://foo/2022-05-21",
                "count": "10"
            }
        },
        ....
    ]

    for m in measurements:
        serialized_measurement = json_metric_serializer_format_to_influxdb_info(m)
        influxdb_utils.write_value_into_influxDB(client=i_client,
                                         tags=serialized_measurement.tags,
                                         fields=serialized_measurement.fields,
                                         measurement=serialized_measurement.measurement,
                                         time=serialized_measurement.time)




def json_metric_serializer_format_to_influxdb_info(meas):
    import collections
    influxdb_info = collections.namedtuple("influxdb_info", ["measurement", "fields", "tags", "time"])

    measurement = str(meas["name"])
    time = meas["fields"].get("time", "")
    fields = {x: meas["fields"][x] for x in meas["fields"] if not x == "time"}
    tags = meas["tags"]

    return influxdb_info(measurement=measurement, fields=fields, tags=tags, time=time)


def write_value_into_influxDB(client, tags, fields, measurement, time):
    import datetime

    if time == '':
        time = datetime.datetime.today().strftime('%Y-%m-%dT%H:%M:%SZ')

    json_body = [
        {
            "measurement": "{measurement}".format(measurement=measurement),
            "tags": tags,
            "fields": fields,
            "time": time
        }
    ]

    client.write_points(json_body)

故障现象

measurements数组可包含上百条条目,但代码运行后仅随机2条条目成功写入name_of_my_measurement测量表,成功写入的条目无明显规律。例如data_path字段存储日期路径,测量数组中包含20条连续日期路径的条目,每次运行脚本后查询InfluxDB,都会有另外随机2条数据落库,现象完全随机。

根本原因

  1. 时间戳精度不足导致同series数据互相覆盖
    InfluxDB判定数据点唯一的依据是measurement名 + 完整tag集 + 时间戳,三个值完全相同的数据点会被判定为同一个点,后写入的会覆盖先写入的。代码中未手动指定时间的点,统一使用秒级精度的写入时刻作为时间戳,而本地循环上百条数据的总耗时通常仅1-2秒,同一秒内写入的所有点tag完全一致(只有customer_code:FOO)、measurement相同、时间戳相同,会被全部覆盖,最终每个秒级时间窗口仅能保留最后1条写入的数据,刚好对应每次仅写入2条左右的现象。
    另外代码用datetime.today()生成本地时间,却硬拼接了UTC的Z后缀,时区不匹配也可能导致时间戳错乱。
  2. 字段类型不匹配导致写入静默失败
    代码中count字段传入的是字符串类型,InfluxDB对字段类型有强约束,同一个measurement下的同名字段一旦确定类型(比如第一次写入是整型),后续传入其他类型的写入请求会直接报错。现有代码没有加任何异常捕获和错误日志,类型不匹配的写入请求会直接静默失败,进一步加剧了随机丢数的情况。

修复方案

  • 提升时间戳精度,不要使用秒级时间戳,可改为微秒/纳秒级UTC时间,例如使用datetime.datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%S.%fZ')生成时间;业务数据最好直接携带对应业务含义的真实时间戳,不要直接用写入时刻作为数据时间。
  • 写入前做字段类型转换,数值类字段比如count要转为int/float类型,不要传字符串。
  • 写入逻辑增加异常捕获,打印写入失败的具体错误和对应数据点,避免静默丢数。
  • 不要循环单条调用write_points,将所有数据点攒成一个列表后批量写入,既提升写入性能,也能减少单条请求的网络波动失败概率。

内容的提问来源于stack exchange,提问作者Niko

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 09:57:16