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

多进程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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 19:29:59