跳转至

脚本开发 / 连接器对象 DFF.CONN / DataKit、DataWay

DataKit、DataWay 连接器操作对象主要提供数据写入方法。

DFF.CONN(...) 参数如下:

参数 类型 必须 / 默认值 说明
connector_id str 必须 连接器 ID
source str None 覆盖连接器 Source
注意不要填写 "mysql" 等采集器名称,以免混淆
timeout int/float 10 当前操作对象的默认 HTTP 请求超时时间,秒
split_size int 100 行协议批量写入时每次请求的数据点数量
参数 类型 必须 / 默认值 说明
connector_id str 必须 连接器 ID
token str None 覆盖连接器 Token
timeout int/float 10 当前操作对象的默认 HTTP 请求超时时间,秒
split_size int 100 行协议批量写入时每次请求的数据点数量
  • 一般性上报数据,请使用 .write_by_category(...).write_by_category_many(...) 方法
  • 一般性执行 DQL 语句,请使用 .query(...) 方法
  • 直接发送 GET 请求,请使用 .get(...) 方法
  • 直接发送 POST 请求,请使用 .post_json(...) 方法
  • 直接发送行协议数据,请使用 .post_line_protocol(...) 方法

本连接器本质上是 HTTP 请求的封装

DataKit 和 DataWay 之间绝大部分接口完全相同。

由于 DataKit、DataWay 接口经常变动,本连接器并不会一对一封装所有的接口

由于不同版本的 DataKit、DataWay 对上报数据可能存在不同的要求或约束,请在阅读相关文档的基础上使用本连接器

详细文档见:

.write_by_category(...)

向 DataKit、DataWay 写入特定类型的数据,参数如下:

参数 类型 必须 / 默认值 说明
category str 必须 数据类型,详见 观测云 文档 / DataKit API
measurement str 必须 指标集名称
tags dict None 标签。键名和键值必须都为字符串
fields dict 必须 指标。键名必须为字符串;值可以为字符串、整数、浮点数、布尔值或元素类型一致的上述类型列表
timestamp int/long/float {当前时间} 时间戳,支持秒/毫秒/微秒/纳秒
headers dict None 请求 Header 参数
timeout int/float None 本次请求超时时间,省略时使用操作对象的默认值

参数 headers 于 3.3.0 新增

示例
1
2
3
tags   = { 'host': 'web-01' }
fields = { 'cpu' : 10 }
status_code, result = datakit.write_by_category(category='metric', measurement='主机监控', tags=tags, fields=fields)

.write_by_category_many(...)

write_by_category(...) 的批量版本,参数如下:

参数 类型 必须 / 默认值 说明
category str 必须 数据类型,详见 观测云 文档 / DataKit API
data list 必须 数据点列表
data[#].measurement str 必须 指标集名称
data[#].tags dict None 标签。键名和键值必须都为字符串
data[#].fields dict 必须 指标。键名必须为字符串;值可以为字符串、整数、浮点数、布尔值或元素类型一致的上述类型列表
data[#].timestamp int/long/float {当前时间} 时间戳,支持秒/毫秒/微秒/纳秒
headers dict None 请求 Header 参数
timeout int/float None 本次请求超时时间,省略时使用操作对象的默认值

参数 headers 于 3.3.0 新增

示例
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
data = [
    {
        'measurement': '主机监控',
        'tags'       : { 'host' : 'web-01' },
        'fields'     : { 'value': 10 }
    },
    {
        'measurement': '主机监控',
        'tags'       : { 'host' : 'web-02' },
        'fields'     : { 'value': 20 }
    }
]
status_code, result = datakit.write_by_category_many(category='metric', data=data)

.write_metric(...) / .write_point(...)

.write_metric(...).write_by_category(category='metric', ...) 等价;.write_point(...) 是旧版兼容别名。

示例
1
status_code, result = datakit.write_metric(measurement='主机监控', tags={'host': 'web-01'}, fields={'cpu': 10})

.write_metric_many(...) / .write_metrics(...) / .write_points(...)

.write_metric_many(...).write_by_category_many(category='metric', ...) 等价;.write_metrics(...).write_points(...) 是旧版兼容别名。

示例
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
data = [
    {
        'measurement': '主机监控',
        'tags'       : { 'host' : 'web-01' },
        'fields'     : { 'value': 10 }
    },
    {
        'measurement': '主机监控',
        'tags'       : { 'host' : 'web-02' },
        'fields'     : { 'value': 20 }
    }
]
status_code, result = datakit.write_metrics(data=data)

.write_logging(...) / .write_logging_many(...)

.write_logging(...).write_by_category(category='logging', ...) 等价;.write_logging_many(...) 是对应的批量版本。

.query(...)

本方法支持 DataKit、DataWay API DQL 查询接口中的参数

详细文档见 观测云 文档 / DataKit API 文档

本方法只是 HTTP 请求的包装

本方法本质上只是向 DataKit、DataWay 发送一个 HTTP 请求,返回内容取决于 DataKit、DataWay 及后端数据源。

如对返回结果有疑问,可以尝试改用 requests 直接向 DataKit、DataWay 发送请求:

使用 requests 调用接口
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
def query():
    domain = '<Domain>'
    token  = '<Token>'

    url = f'https://{domain}/v1/query/raw?token={token}'
    body = {
        'queries': [
            {
                # DQL 语句
                'query': 'M::`cpu`:(`load5s`) BY `host`',

                # 最近 1 小时
                'time_range': [
                    _DFF_TRIGGER_TIME_MS - 3600 * 1000,
                    _DFF_TRIGGER_TIME_MS,
                ],
            }
        ],
        'token': token
    }

    resp = requests.post(url, json=body)
    print(resp.status_code)
    print(resp.text)

通过 DataKit、DataWay 执行 DQL 语句,参数如下:

参数 类型 必须 / 默认值 说明
dql str 必须 DQL 语句
dict_output bool False 是否自动转换数据为 dict
raw bool False 是否返回原始响应。开启后 dict_output 参数无效。
all_series bool False 是否自动通过 slimitsoffset 翻页以获取全部时间线。
token str None 仅 DataKit 使用的工作空间 Token;DataWay 应在获取操作对象时设置 Token
timeout int/float None 本次请求超时时间,省略时使用操作对象的默认值
{DataKit、DataWay 原生参数} - - 透传至 queries[0].{DataKit、DataWay 原生参数}

启用 all_series 时,每页固定查询 500 条时间线:指标查询最多请求 20 页,其他查询最多请求 5 页。

DataWay 执行查询前必须具有 Token。应在连接器配置或 DFF.CONN(..., token='...') 中设置,不要向 .query(...) 重复传入 token

示例
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
import time
import json

@DFF.API('Run DQL via DataKit')
def run_dql_via_datakit():
    datakit = DFF.CONN('datakit')

    # 使用 DataKit 原生参数 `time_range`,限制最近 1 小时数据
    time_range = [
        int(time.time() - 3600) * 1000,
        int(time.time()) * 1000,
    ]

    # 查询并以 dict 形式返回数据
    status_code, result = datakit.query(dql='O::HOST:(host,load,create_time)', dict_output=True, time_range=time_range)
    print(json.dumps(result))
输出示例
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
{
  "series": [
    [
      {
        "time": 1622463105293,
        "host": "iZbp152ke14timzud0du15Z",
        "load": 2.18,
        "create_time": 1622429576363,
        "tags": {}
      },
      {
        "time": 1622462905921,
        "host": "ubuntu18-base",
        "load": 0.08,
        "create_time": 1622268259114,
        "tags": {}
      },
      {
        "time": 1622461264175,
        "host": "shenrongMacBook.local",
        "load": 2.395508,
        "create_time": 1622427320834,
        "tags": {}
      }
    ]
  ]
}
示例
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
import time
import json

@DFF.API('Run DQL via DataKit')
def run_dql_via_datakit():
    datakit = DFF.CONN('datakit')

    # 添加 raw 参数,获取 DQL 查询原始值
    time_range = [
        int(time.time() - 3600) * 1000,
        int(time.time()) * 1000,
    ]

    # 查询并以 DataKit 原始返回值格式返回数据
    status_code, result = datakit.query(dql='O::HOST:(host,load,create_time)', raw=True, time_range=time_range)
    print(json.dumps(result, indent=2))
输出示例
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
{
  "content": [
    {
      "series": [
        {
          "name": "HOST",
          "columns": [
            "time",
            "host",
            "load",
            "create_time"
          ],
          "values": [
            [
              1622463165152,
              "iZbp152ke14timzud0du15Z",
              1.92,
              1622429576363
            ],
            [
              1622462905921,
              "ubuntu18-base",
              0.08,
              1622268259114
            ],
            [
              1622461264175,
              "shenrongMacBook.local",
              2.395508,
              1622427320834
            ]
          ]
        }
      ],
      "cost": "1ms",
      "total_hits": 3
    }
  ]
}

.get(...)

本方法为通用处理方法

具体参数格式、内容等请参考 观测云 文档 / DataKit API

向 DataKit、DataWay 发送一个 GET 请求,参数如下:

参数 类型 必须 / 默认值 说明
path str 必须 请求路径
query dict None 请求 URL 参数
headers dict None 请求 Header 参数
timeout int/float None 本次请求超时时间,省略时使用操作对象的默认值

返回 (status_code, result)。响应体可解析为 JSON 时,result 为对应对象,否则为文本或原始内容。

.post_json(...)

本方法为通用处理方法

具体参数格式、内容等请参考 观测云 文档 / DataKit API

向 DataKit、DataWay 以 JSON 格式发送一个 POST 请求,参数如下:

参数 类型 必须 / 默认值 说明
path str 必须 请求路径
json_obj dict/list 必须 需要发送的 JSON 对象
query dict None 请求 URL 参数
headers dict None 请求 Header 参数
timeout int/float None 本次请求超时时间,省略时使用操作对象的默认值

参数 path 于 1.6.8 版本调整为第一个参数

返回 (status_code, result)

.post_line_protocol(...)

本方法为通用处理方法

具体参数格式、内容等请参考 观测云 文档 / DataKit API

向 DataKit、DataWay 以行协议格式发送一个 POST 请求,参数如下:

参数 类型 必须 / 默认值 说明
path str 必须 请求路径
points dict/list 必须 单个数据点或数据点列表
points[#].measurement str 必须 指标集名称
points[#].tags dict None 标签。键名和键值必须都为字符串
points[#].fields dict 必须 指标。键名必须为字符串;值可以为字符串、整数、浮点数、布尔值或元素类型一致的上述类型列表
points[#].timestamp int/long/float {当前时间} 时间戳,支持秒/毫秒/微秒/纳秒
query dict None 请求 URL 参数
headers dict None 请求 Header 参数
timeout int/float None 本次请求超时时间,省略时使用操作对象的默认值

参数 path 于 1.6.8 版本调整为第一个参数

批量数据按 split_size 分片发送。本方法返回最后一次分片请求的 (status_code, result);任一分片请求失败时抛出异常并停止后续发送。