使用 Python 查询与写入 Loki:HTTP API 客户端实战指南
【免费下载链接】lokiLike Prometheus, but for logs.项目地址: https://gitcode.com/GitHub_Trending/lok/loki
Loki 的 HTTP API(/loki/api/v1/*)是它与外部系统交互的通用接口,本文基于 Loki 官方 Python 客户端示例,完整演示如何用requests与httpx完成日志范围查询、即时指标查询、日志推送、标签枚举等高频操作,并覆盖多租户认证、Grafana Cloud 接入、TLS 校验与错误重试等生产环境细节。读完本文,你将掌握一套可直接复制运行的 Python 代码骨架,并能结合 Loki HTTP API 参考 深入每个端点的全部参数与响应格式。
前置条件:安装 HTTP 客户端库
先安装同步客户端 requests:
pip install requests如果你需要异步能力,httpx 提供了几乎一致的 API:
pip install httpx两种库在本文示例中的调用方式完全对称,唯一的细微差别是httpx的verify参数还额外接受ssl.SSLContext对象。下述所有函数都统一接收headers、auth、verify三个可选参数,方便你在不同部署环境下复用同一套代码。
认证:从本地实例到多租户与 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>"), )User与URL两个值都可以在 Grafana Cloud Stack 的 Loki 日志服务详情页中找到。注意auth元组同时兼容requests与httpx的签名,代码无需改动。
自签名证书:控制 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在一个时间区间内查询日志,是最常用的查询操作(也用于返回日志行的时间范围查询)。它由querier、query-frontend、read和all组件暴露,具体见 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)解析,默认值为backward(defaultDirection常量,见 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-Type为application/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/labels与GET /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 状态码,常见错误包括:
| 状态码 | 含义 | 常见原因 |
|---|---|---|
| 400 | Bad Request | LogQL 语法无效 |
| 429 | Too Many Requests | 触发速率限制 |
| 5xx | Server Error | Loki 不可用或过载 |
生产客户端应通过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_range,rate()/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),仅供参考