news 2026/9/12 18:47:50

使用 Python 查询与写入 Loki:HTTP API 客户端实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
使用 Python 查询与写入 Loki:HTTP API 客户端实战指南

使用 Python 查询与写入 Loki:HTTP API 客户端实战指南

【免费下载链接】lokiLike Prometheus, but for logs.项目地址: https://gitcode.com/GitHub_Trending/lok/loki

Loki 的 HTTP API(/loki/api/v1/*)是它与外部系统交互的通用接口,本文基于 Loki 官方 Python 客户端示例,完整演示如何用requestshttpx完成日志范围查询、即时指标查询、日志推送、标签枚举等高频操作,并覆盖多租户认证、Grafana Cloud 接入、TLS 校验与错误重试等生产环境细节。读完本文,你将掌握一套可直接复制运行的 Python 代码骨架,并能结合 Loki HTTP API 参考 深入每个端点的全部参数与响应格式。

前置条件:安装 HTTP 客户端库

先安装同步客户端 requests:

pip install requests

如果你需要异步能力,httpx 提供了几乎一致的 API:

pip install httpx

两种库在本文示例中的调用方式完全对称,唯一的细微差别是httpxverify参数还额外接受ssl.SSLContext对象。下述所有函数都统一接收headersauthverify三个可选参数,方便你在不同部署环境下复用同一套代码。

认证:从本地实例到多租户与 Grafana Cloud

本文示例默认连接无认证的本地 Loki 实例(如http://localhost:3100)。如果你的集群启用了多租户,或接入了 Grafana Cloud,需要按下述方式补充认证信息。

多租户模式:X-Scope-OrgID请求头

当集群开启多租户隔离后,每个请求都必须携带租户 ID。直接在请求头中注入X-Scope-OrgID

headers = {"X-Scope-OrgID": "my-tenant"} resp = requests.get(url, params=params, headers=headers)

套用到本文定义的函数上:

results = query_range( url="http://localhost:3100", query='{job="varlogs"}', start=datetime.now() - timedelta(hours=1), end=datetime.now(), headers={"X-Scope-OrgID": "my-tenant"}, )

要跨多个租户查询,用管道符(|)分隔租户名:

headers = {"X-Scope-OrgID": "tenant1|tenant2|tenant3"}

Grafana Cloud:Basic Auth

对 Grafana Cloud 的 Loki 服务,使用 Basic Auth 传入你的 Grafana Cloud 用户和 API Token:

resp = requests.get(url, params=params, auth=("<user>", "<API_TOKEN>"))

套用到函数上:

results = query_range( url="https://logs-prod-us-central1.grafana.net", query='{job="varlogs"}', start=datetime.now() - timedelta(hours=1), end=datetime.now(), auth=("<user>", "<API_TOKEN>"), )

UserURL两个值都可以在 Grafana Cloud Stack 的 Loki 日志服务详情页中找到。注意auth元组同时兼容requestshttpx的签名,代码无需改动。

自签名证书:控制 TLS 校验

本地或内网 Loki 常使用自签名 TLS 证书,此时可临时跳过证书校验:

resp = requests.get(url, params=params, verify=False)

生产环境不建议关闭校验,改为传入 CA 证书包路径:

resp = requests.get(url, params=params, verify="/path/to/ca-bundle.crt")

套用到函数上:

results = query_range( url="https://loki.internal:3100", query='{job="varlogs"}', start=datetime.now() - timedelta(hours=1), end=datetime.now(), verify="/path/to/ca-bundle.crt", )

范围查询:GET /loki/api/v1/query_range

/loki/api/v1/query_range一个时间区间内查询日志,是最常用的查询操作(也用于返回日志行的时间范围查询)。它由querierquery-frontendreadall组件暴露,具体见 Loki HTTP API 参考。

使用 requests

import requests from datetime import datetime, timedelta def query_range( url: str, query: str, start: datetime, end: datetime, limit: int = 1000, headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle ) -> list: """Query Loki for log entries within a time range.""" resp = requests.get( f"{url}/loki/api/v1/query_range", params={ "query": query, "start": str(int(start.timestamp() * 1e9)), # nanoseconds "end": str(int(end.timestamp() * 1e9)), "limit": limit, "direction": "backward", }, headers=headers, auth=auth, verify=verify, ) resp.raise_for_status() return resp.json()["data"]["result"] results = query_range( url="http://localhost:3100", query='{job="varlogs"} |= "error"', start=datetime.now() - timedelta(hours=1), end=datetime.now(), ) for stream in results: print(f"Labels: {stream['stream']}") for ts, line in stream["values"]: print(f" {datetime.fromtimestamp(int(ts) / 1e9)}: {line}")

使用 httpx

import httpx from datetime import datetime, timedelta def query_range( url: str, query: str, start: datetime, end: datetime, limit: int = 1000, headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle; httpx also accepts ssl.SSLContext ) -> list: """Query Loki for log entries within a time range.""" resp = httpx.get( f"{url}/loki/api/v1/query_range", params={ "query": query, "start": str(int(start.timestamp() * 1e9)), # nanoseconds "end": str(int(end.timestamp() * 1e9)), "limit": limit, "direction": "backward", }, headers=headers, auth=auth, verify=verify, ) resp.raise_for_status() return resp.json()["data"]["result"] results = query_range( url="http://localhost:3100", query='{job="varlogs"} |= "error"', start=datetime.now() - timedelta(hours=1), end=datetime.now(), ) for stream in results: print(f"Labels: {stream['stream']}") for ts, line in stream["values"]: print(f" {datetime.fromtimestamp(int(ts) / 1e9)}: {line}")

参数背后的实现细节

从 pkg/loghttp/params.go 的源码可以确认这些参数在服务端的解析规则:

  • start/end:由parseTimestamp(params.go#L175-L198)解析,支持纳秒整数时间戳、含小数点的浮点时间戳,甚至是RFC3339Nano格式的字符串;位数不超过 10 位时按秒处理。官方示例统一用纳秒字符串,最稳妥。
  • direction:由parseDirection(params.go#L202-L212)解析,默认值为backwarddefaultDirection常量,见 params.go#L24),即最新日志在前。
  • limit:默认值在 params.go#L21 定义为defaultQueryLimit = 100,与文档“Common problems”一节描述的默认 100 条一致。本文示例显式传limit=1000覆盖该默认值。
  • since:除start/end外,服务端还支持since参数(如since=2h),缺省start时会按end减去since计算,默认 1 小时(defaultSince,见 params.go#L23)。
  • step:范围查询的步长可省略,服务端会按时间跨度动态计算默认步长max(floor((end-start)/250), 1)秒(defaultQueryRangeStep,见 params.go#L149-L151)。

响应体结构定义在 pkg/loghttp/query.go 的QueryResponse(query.go#L55-L60),包含status、可选的warnings以及data字段;data.result是一个数组,日志查询中每个元素含stream(标签集合)与values[纳秒时间戳, 日志行]列表),这正是示例中遍历结构的依据。

即时查询:GET /loki/api/v1/query

/loki/api/v1/query单个时间点(默认当前时刻)评估查询,适用于rate()count_over_time()bytes_over_time()等聚合构成的即时指标查询。注意:返回日志行的流选择器(如{job="myapp"})不支持即时查询,必须改用query_range

import requests from datetime import datetime def query_instant( url: str, query: str, headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle ) -> list: """Run an instant metric query against Loki.""" resp = requests.get( f"{url}/loki/api/v1/query", params={ "query": query, "time": str(int(datetime.now().timestamp() * 1e9)), }, headers=headers, auth=auth, verify=verify, ) resp.raise_for_status() return resp.json()["data"]["result"] results = query_instant( url="http://localhost:3100", query='sum(rate({job="varlogs"}[10m])) by (level)', ) for entry in results: print(f"{entry['metric']}: {entry['value'][1]}")

time参数同样由parseTimestamp解析(见 params.go#L80-L82),缺省取当前时间,因此最小可用的即时查询甚至可以只传query一个参数。指标查询的响应中每个结果元素含metric(分组标签)与value[时间戳, 值]二元组),与示例中的打印逻辑一一对应。

推送日志:POST /loki/api/v1/push

/loki/api/v1/push是向 Loki 写入日志的端点,在微服务模式下由distributor暴露。默认的 POST body 是 Snappy 压缩的 Protobuf 消息,但当Content-Typeapplication/json时可以直接发送 JSON,格式为:

{ "streams": [ { "stream": { "label": "value" }, "values": [ [ "<unix epoch in nanoseconds>", "<log line>" ], [ "<unix epoch in nanoseconds>", "<log line>" ] ] } ] }

对应到 Python 实现:

import json import time import requests def push_logs( url: str, labels: dict[str, str], entries: list[tuple[str, str]], headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle ) -> None: """Push log entries to Loki. Args: url: Loki base URL. labels: Stream labels, for example {"job": "myapp", "env": "dev"}. entries: List of (timestamp_ns, log_line) tuples. Use str(int(time.time() * 1e9)) to get a nanosecond timestamp. """ payload = { "streams": [ { "stream": labels, "values": [list(e) for e in entries], } ] } req_headers = {**(headers or {}), "Content-Type": "application/json"} resp = requests.post( f"{url}/loki/api/v1/push", headers=req_headers, data=json.dumps(payload), auth=auth, verify=verify, ) resp.raise_for_status() now_ns = str(int(time.time() * 1e9)) push_logs( url="http://localhost:3100", labels={"job": "myapp", "env": "dev"}, entries=[ (now_ns, "application started"), (now_ns, "listening on port 8080"), ], )

JSON Push 的几个关键约束

结合 HTTP API 文档的 Ingest 章节:

  • 时间戳必须传字符串而非数字,否则端点会返回 400 错误;示例用str(int(time.time() * 1e9))构造纳秒字符串。
  • 每行日志必须是[纳秒时间戳, 日志行]二元组,写入时转成 JSON 数组(list(e))。
  • 可以用Content-Encoding: gzip请求头发送 gzip 压缩的 JSON 体,降低大流量推送的网络开销。
  • 可选地在每个日志行的数组末尾追加一个结构化元数据JSON 对象(键和值都必须是字符串,不允许嵌套),例如["<纳秒时间戳>", "<日志行>", {"trace_id": "0242ac120002", "user_id": "superUser123"}],配合结构化元数据文档使用。
  • 服务端对 JSON 请求体的解析在 pkg/loghttp/query.go 的PushRequest/LogProtoStream中实现(query.go#L100-L120),JSON 流会被转换为内部logproto.Stream的标签串格式,因此“stream”里的标签键值必须是合法标签。

查询标签与标签值:GET /loki/api/v1/labelsGET /loki/api/v1/label/<name>/values

/loki/api/v1/labels返回所有已知标签名列表;/loki/api/v1/label/<name>/values返回某个标签的全部取值。两者常用于驱动探索式分析或构建日志筛选 UI。

import requests def get_labels( url: str, headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle ) -> list[str]: """List all known label names.""" resp = requests.get( f"{url}/loki/api/v1/labels", headers=headers, auth=auth, verify=verify, ) resp.raise_for_status() return resp.json()["data"] def get_label_values( url: str, label: str, headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle ) -> list[str]: """List values for a specific label.""" resp = requests.get( f"{url}/loki/api/v1/label/{label}/values", headers=headers, auth=auth, verify=verify, ) resp.raise_for_status() return resp.json()["data"] labels = get_labels("http://localhost:3100") print(f"Labels: {labels}") for label in labels: values = get_label_values("http://localhost:3100", label) print(f" {label}: {values}")

注意响应中的data字段是一个字符串数组(标签名或标签值列表),与查询接口的data.result结构不同,直接索引返回即可。

错误处理:HTTP 状态码与指数退避重试

Loki 返回标准的 HTTP 状态码,常见错误包括:

状态码含义常见原因
400Bad RequestLogQL 语法无效
429Too Many Requests触发速率限制
5xxServer ErrorLoki 不可用或过载

生产客户端应通过raise_for_status()捕获错误,并结合响应体排查详情。针对 429 这类可恢复错误,可加入指数退避重试:

import time import requests def query_with_retry( url: str, query: str, max_retries: int = 3, backoff: float = 1.0, headers: dict[str, str] | None = None, auth: tuple[str, str] | None = None, verify: bool | str = True, # False to skip TLS, or path to CA bundle ) -> dict: """Query Loki with simple retry logic for rate limits.""" for attempt in range(max_retries): resp = requests.get( f"{url}/loki/api/v1/query", params={"query": query}, headers=headers, auth=auth, verify=verify, ) if resp.status_code == 429: wait = backoff * (2 ** attempt) print(f"Rate limited, retrying in {wait}s...") time.sleep(wait) continue resp.raise_for_status() return resp.json() raise Exception(f"Query failed after {max_retries} retries")

重试等待时间按backoff * 2^attempt指数增长(1s、2s、4s……),在不超过max_retries的前提下自动降频,避免在限流期间继续冲击服务端。关于请求级与租户级限流配置,可进一步参考请求校验与速率限制文档。

常见问题速查

  • 时间戳必须是纳秒。Loki 期望 Unix 纳秒时间戳,而不是秒或毫秒。用time.time() * 1e9换算后转成字符串;查询端别忘了把响应中的纳秒除以1e9再转datetime展示。
  • 至少需要一个标签匹配器。没有流选择器就无法查询。{job="myapp"}合法,空选择器不合法。
  • direction参数决定结果排序。backward(默认)先返回最新条目,forward先返回最旧条目。范围查询时分页依赖这个语义,见下一条。
  • 即时查询端点只支持指标查询。/query上执行{job="myapp"}这类日志流选择器会返回 400,日志查询请走/query_rangerate()/count_over_time()等聚合走/query
  • limit控制结果规模并实现分页。查询类接口的limit默认值为 100(见 params.go#L21)。对于大时间范围,可调高limit,或根据返回的最后一条时间戳把start前移后继续拉取,实现游标式分页;query_range还支持按step控制指标计算的分辨率。

延伸阅读

  • 完整的端点清单、全部查询参数与响应格式,参见 Loki HTTP API 参考。
  • 多租户机制的配置与原理,参见多租户文档。
  • 本文示例查询中使用的 LogQL 语法与聚合函数,参见 LogQL 查询文档。
  • 服务端对query_range/query/labels各参数的解析实现,可深入 pkg/loghttp/params.go 与 pkg/loghttp/query.go 阅读源码。

【免费下载链接】lokiLike Prometheus, but for logs.项目地址: https://gitcode.com/GitHub_Trending/lok/loki

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 18:40:48

基于Neo4j的医疗知识图谱问答系统:从实体识别到Cypher查询

简介&#xff1a;基于Python的知识图谱医疗领域问答系统项目&#xff0c;面向Python学习者和知识图谱入门者&#xff0c;提供一套可完整运行的医疗领域问答系统源码与配套数据&#xff0c;适合作为期末大作业或课程设计参考。整个项目针对医疗知识结构化表示与自动问答实现进行…

作者头像 李华
网站建设 2026/9/12 18:40:28

SAP S/4HANA信用管理核心接口视图I_CreditManagementBP详解

干过 SAP Credit Management 项目的朋友都有一种感觉&#xff1a;信用管理这个模块&#xff0c;功能强大&#xff0c;但数据结构的复杂程度也是出了名的。尤其是当你想快速回答一个最简单的业务问题——“这个业务伙伴到底能不能给他放额度、放多少”——的时候&#xff0c;往往…

作者头像 李华
网站建设 2026/9/12 18:38:20

ML-KWS-for-MCU源码评测:Cortex-M语音唤醒的完整AI链路

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/12 18:29:51

基于ADS1278 EKIT的8通道同步数据采集:时序、代码与验证

简介&#xff1a;这是围绕 TI ADS1278 高精度 24 位 Σ-Δ ADC 在 TI DSP 平台上的参考工程代码包&#xff0c;面向数据采集、传感器接口与信号调理等嵌入式开发者&#xff0c;帮助快速掌握芯片驱动与系统集成方法。资源共 44 个文件&#xff0c;压缩包约 445KB&#xff0c;包含…

作者头像 李华