多进程Hypercorn部署下Quart应用的Prometheus指标聚合方案咨询
Hypercorn多进程Quart应用Prometheus指标汇总方案
问题场景
生产环境用Hypercorn运行Quart应用,配置8个工作进程,基于aioprometheus库统计请求延迟、吞吐量等指标,通过/myapp/metrics端点暴露给Prometheus采集。但由于请求会被路由到任意一个工作进程,每次采集只能拿到当前进程的指标数据,无法汇总所有8个进程的统计结果(比如进程A统计6次事件、进程B统计7次,采集时只能拿到6或7,得不到总数13)。
下面针对Hypercorn的特性提供专属解决方案,无需依赖外部公共数据源:
方案1:主进程转发采集请求到所有工作进程汇总
Hypercorn的主进程负责管理所有子工作进程,我们可以利用这一点,让主进程接收Prometheus的采集请求后,主动拉取每个工作进程的内部指标,再合并汇总返回。
实现步骤:
- 每个工作进程启动时,绑定一个仅本地可访问的内部端口,暴露私有metrics端点(避免对外暴露)
- 主进程部署一个轻量ASGI服务,处理对外的
/myapp/metrics请求:遍历所有工作进程,逐个请求它们的内部metrics接口,拿到数据后按指标类型合并(计数器求和、直方图合并分桶等) - 配置Hypercorn让主进程监听对外端口,工作进程仅监听本地内部端口
代码示例:
工作进程代码(worker.py):
import os from quart import Quart from aioprometheus import Counter, MetricsMiddleware app = Quart(__name__) app.asgi_app = MetricsMiddleware(app.asgi_app) # 定义业务指标 request_counter = Counter("request_total", "Total requests processed") request_duration = Histogram("request_duration_seconds", "Request processing duration") @app.route("/") async def index(): request_counter.inc({"endpoint": "/"}) return "Hello from worker" # 分配本地专属端口,避免冲突 if __name__ == "__main__": worker_port = 8000 + int(os.getpid() % 8) # 对应8个工作进程的端口段 await app.run(host="127.0.0.1", port=worker_port)
主进程汇总代码(main.py):
import asyncio import httpx from hypercorn.config import Config from hypercorn.asyncio import serve from quart import Quart, Response main_app = Quart(__name__) # 预定义工作进程的内部端口(也可以从Hypercorn进程列表动态获取) WORKER_PORTS = [8000 + i for i in range(8)] async def fetch_worker_metrics(port): """拉取单个工作进程的指标数据""" async with httpx.AsyncClient() as client: resp = await client.get(f"http://127.0.0.1:{port}/metrics") return resp.text def merge_metrics(metrics_list): """合并多个进程的Prometheus格式指标""" merged = {} for metrics_text in metrics_list: for line in metrics_text.split("\n"): line = line.strip() if not line or line.startswith("#"): continue # 拆分指标名、标签、数值(简化处理,实际需兼容复杂标签) metric_part, value_part = line.rsplit(" ", 1) value = float(value_part) if metric_part not in merged: merged[metric_part] = 0.0 merged[metric_part] += value # 重新生成Prometheus格式文本 result = [] # 先保留注释头(取第一个进程的注释) result.extend([line for line in metrics_list[0].split("\n") if line.startswith("#")]) for metric, val in merged.items(): result.append(f"{metric} {val}") return "\n".join(result) @main_app.route("/myapp/metrics") async def aggregated_metrics(): """对外暴露的汇总指标端点""" tasks = [fetch_worker_metrics(port) for port in WORKER_PORTS] all_metrics = await asyncio.gather(*tasks) merged_text = merge_metrics(all_metrics) return Response(merged_text, content_type="text/plain; version=0.0.4") # 启动Hypercorn工作进程+主进程汇总服务 async def run(): # 配置Hypercorn启动8个工作进程 worker_config = Config() worker_config.workers = 8 worker_config.application_path = "worker:app" worker_config.bind = [f"127.0.0.1:{port}" for port in WORKER_PORTS] # 同时启动主进程对外服务和工作进程 task1 = asyncio.create_task(serve(main_app, Config(bind=["0.0.0.0:8080"]))) task2 = asyncio.create_task(serve(worker_config.application, worker_config)) await asyncio.gather(task1, task2) if __name__ == "__main__": asyncio.run(run())
方案2:利用Hypercorn的fork模型共享内存存储指标
Hypercorn的工作进程是通过fork主进程创建的,子进程会继承主进程的内存空间。我们可以用Python的共享内存机制,让所有进程共享同一个指标存储实例,避免数据分散。
实现步骤:
- 在主进程fork工作进程之前,创建共享内存区域,初始化aioprometheus指标并绑定到共享内存
- 自定义指标类,将更新和读取操作指向共享内存,确保所有进程操作同一份数据
- 注意处理并发更新的原子性,避免竞态条件
代码示例:
from multiprocessing import shared_memory import struct from aioprometheus import Counter, Registry from quart import Quart, Response from hypercorn.config import Config from hypercorn.asyncio import serve # 主进程创建共享内存(需在fork前执行) shm = shared_memory.SharedMemory(create=True, size=1024) # 按需调整大小 class SharedCounter(Counter): """基于共享内存的计数器,支持多进程共享""" def __init__(self, name, description): super().__init__(name, description) # 初始化共享内存值为0 shm.buf[:8] = struct.pack("Q", 0) def inc(self, labels=None, value=1): # 原子递增(简化实现,生产环境建议用锁或更安全的原子操作) while True: current = struct.unpack("Q", shm.buf[:8])[0] new_val = current + value # 尝试写入,避免竞态 shm.buf[:8] = struct.pack("Q", new_val) if struct.unpack("Q", shm.buf[:8])[0] == new_val: break def get(self): return struct.unpack("Q", shm.buf[:8])[0] # 主进程初始化指标 request_counter = SharedCounter("request_total", "Total requests processed") registry = Registry() registry.register(request_counter) app = Quart(__name__) @app.route("/") async def index(): request_counter.inc() return "Hello from shared memory worker" @app.route("/myapp/metrics") async def metrics(): # 生成Prometheus格式指标 metrics_text = ( "# HELP request_total Total requests processed\n" "# TYPE request_total counter\n" f"request_total {request_counter.get()}" ) return Response(metrics_text, content_type="text/plain; version=0.0.4") if __name__ == "__main__": config = Config() config.workers = 8 config.bind = ["0.0.0.0:8080"] # 启动时fork工作进程,子进程继承共享内存 asyncio.run(serve(app, config))
方案对比
- 主进程转发方案:实现简单,无需修改指标逻辑,兼容性好,但需要额外的内部请求开销,适合指标量不大的场景
- 共享内存方案:性能更高,无额外网络开销,但需要处理并发竞态,对指标类型的兼容性有限(比如直方图的分桶共享需要更复杂的实现)
内容的提问来源于stack exchange,提问作者shshnk
相关产品推荐
相关产品推荐

