Langflow Telemetry 写入压测:用 stress_telemetry_writes.py 复现并验证 DB 连接池耗尽问题

【免费下载链接】langflow Langflow is a powerful tool for building and deploying AI-powered agents and workflows. 【免费下载链接】langflow 项目地址: https://gitcode.com/GitHub_Trending/la/langflow

本文基于 Langflow 仓库中的 telemetry 写入压测工具(src/backend/tests/stress/)展开。该工具是一个手动压测 harness,用于复现 "transactions + vertex_builds 写入耗尽数据库连接池" 这一经典故障模式,并在 SQLite 与 Postgres 两种后端上验证修复效果。读完本文,你将掌握:如何在不依赖 OpenAI key、不运行真实 flow 的前提下对 Langflow 遥测写入链路施压;如何通过环境变量复现旧版直接写入(legacy direct-write)路径下的 database is locked / 池超时故障;以及如何解读压测输出中的 error_total、enq_tx/enq_vb 计数与行数对账结果。

一、压测目标:绕过 HTTP 层,直击遥测写入路径

Langflow 在重负载 flow 执行时,每个顶点(vertex)完成后都会通过 BackgroundTasks 工作协程异步写入两类遥测行:

  • transaction 表:由 transaction_service.log_transaction() 记录组件执行的输入、输出与状态;
  • vertex_build 表:由 lfx.graph.utils.log_vertex_build() 记录每个顶点的构建产物。

这两类写入默认共享 FastAPI 请求处理所使用的数据库连接池。当并发足够高时,遥测写入会"吃掉"请求池,导致 SQLite 报 database is locked、Postgres 报池超时。stress_telemetry_writes.py 的设计目标是绕过 FastAPI / locust / flow 执行层,从大量并发协程中直接调用上述两个写入函数——这正好模拟了 per-vertex BackgroundTasks 工作器在重负载下的行为,同时不需要 OpenAI key,也不需要真实的 flow(见 工具文档)。

从 harness 源码 可以看到 worker 的实现要点:

tx_service = get_transaction_service()
# 每轮迭代依次调用两个被压函数
await tx_service.log_transaction(
    flow_id=flow_id,
    vertex_id=f"vertex-{worker_id}-{i % 5}",
    inputs=make_inputs(i), outputs=make_outputs(i),
    status="success", target_id=f"target-{worker_id}-{i % 5}",
)
await log_vertex_build(
    flow_id=flow_id, vertex_id=f"vertex-{worker_id}-{i % 5}",
    valid=True, params=str(make_inputs(i)), data=make_outputs(i), artifacts=None,
)
await asyncio.sleep(random.uniform(0.0, 0.01))
  • 每个 worker 是一个 asyncio 协程,循环调用写入函数直到 --seconds 到期;
  • vertex_id 按 i % 5 复用 5 个虚拟顶点,避免构造无意义的海量 id;
  • 每次写入之间随机休眠 0~10ms,形成近似泊松到达的写负载;
  • 每次调用独立 try/except,按异常类型(如 tx:OperationalError、vb:TimeoutError)计入计数器 counters,并保留前 10 条错误样本用于事后分析。

harness 在导入 langflow 之前先设置环境变量,确保 settings 加载到正确的 DB 与开关(源码 L31-L37):

环境变量作用
DB_URL目标数据库 DSN,缺省为 sqlite:///./stress.db;会注入 LANGFLOW_DATABASE_URL
LANGFLOW_TRANSACTIONS_STORAGE_ENABLED / LANGFLOW_VERTEX_BUILDS_STORAGE_ENABLED强制开启两类遥测存储
LANGFLOW_AUTO_LOGIN设为 true,避免服务初始化阻塞在用户配置路径上
DRAIN_GRACE_Sworker 结束后等待 writer 排空的宽限秒数,缺省 5 秒(源码 L181)

二、三种运行方式(完整继承自工具文档)

2.1 对 SQLite 压测

uv run python src/backend/tests/stress/stress_telemetry_writes.py \
    --concurrency 200 --seconds 15

SQLite 模式下会在工作目录生成 stress.db,是验证 "database is locked" 是否被修复的最快路径。

2.2 对 Postgres 压测

# Start a throwaway Postgres
docker run --rm -d --name langflow-pg-test -p 55432:5432 \
    -e POSTGRES_PASSWORD=langflow -e POSTGRES_USER=langflow \
    -e POSTGRES_DB=langflow postgres:16

export DB_URL="postgresql+psycopg://langflow:langflow@localhost:55432/langflow"
# (or any standard postgres DSN, e.g. "$LANGFLOW_DATABASE_URL" if set)
uv run python src/backend/tests/stress/stress_telemetry_writes.py \
    --concurrency 200 --seconds 15

注意 DSN 使用 postgresql+psycopg 方言,与 Langflow 异步 SQLAlchemy 引擎一致;端口 55432 是刻意避开本机已有 Postgres 的一次性端口。

2.3 关闭 writer,复现旧版故障模式

LANGFLOW_TELEMETRY_WRITER_ENABLED=false \
LANGFLOW_DB_CONNECTION_SETTINGS='{"pool_size":5,"max_overflow":5,"pool_timeout":3}' \
    uv run python src/backend/tests/stress/stress_telemetry_writes.py \
        --concurrency 500 --seconds 15

这一组合做两件事:

  1. LANGFLOW_TELEMETRY_WRITER_ENABLED=false 强制所有写入走 legacy 直接写入路径(每条 INSERT 都占用一个请求池连接,见 TransactionService.log_transaction);
  2. LANGFLOW_DB_CONNECTION_SETTINGS 把连接池缩到极小:pool_size=5, max_overflow=5, pool_timeout=3。作为对照,数据库设置组的默认值 是 pool_size=20, max_overflow=30, pool_timeout=30——在 500 并发协程、每次写入都要持有一个连接完成事务的情况下,5+5 的连接、3 秒的等待上限必然被击穿,从而稳定复现池超时(Postgres)或锁竞争(SQLite)。

harness 命令行参数(argparse 定义):

参数缺省值说明
--concurrency200并发 worker 协程数
--seconds20生产阶段(producer phase)持续时间(秒)

进程退出码即压测结论:error_total > 0 时返回 1,否则返回 0(源码 L207-L209),可直接用于 CI 或脚本化回归。

三、输出解读:三个对账点

工具文档给出了判断依据(What to look for),结合 harness 打印逻辑逐条对应:

  1. error_total=0 —— writer 启用时,既不应出现 SQLite "database is locked",也不应出现 Postgres 池超时异常。error_total 是所有非 *_ok 计数(即按异常类型聚合的 tx:*、vb:* 与 worker 级异常)之和。

  2. enq_tx / enq_vb 计数与生产者成功数对齐 —— harness 会在 "workers-done" 与 "post-drain" 两个时间点打印 writer 状态(源码 L133-L140):

    [run]  writer @ post-drain: enq_tx=... enq_vb=... flushed_rows=...
          failed_batches=... dropped_tx=... dropped_vb=...
          tx_buffer=... vb_buffer=...
    

    enqueued_transactions / enqueued_vertex_builds 应等于 counters 中 transactions_ok + vertex_builds_ok;failed_batches=0、dropped_* 为 0 表示无批量写入失败与队列溢出丢弃。

  3. 排空宽限后,DB 行增量应等于入队量(扣除 retention 修剪) —— harness 在跑前、跑后各执行一次 SELECT count(*) FROM "transaction" / "vertex_build"(count_rows),并打印 delta。由于 retention sweeper 可能已在保留上限内修剪旧行(见下文 max_transactions_to_keep 等配置),对账是"模 retention"的近似相等,而非严格相等。

四、底层原理:TelemetryWriterService 如何让遥测写入不再耗尽请求池

理解 harness 压的是什么,需要看它保护的对象:TelemetryWriterService。其设计在源码 docstring 中已完整声明:把 transaction 与 vertex_build 的写入解耦出请求处理连接池。

4.1 写入链路:内存队列 + 专用小连接池

  • 生产端是 O(1) 内存入队:enqueue_transaction / enqueue_vertex_build 只是把行 append 到 collections.deque,不获取任何 DB 连接(源码 L261-L280)。
  • 单后台 writer 任务批量落库:_run_writer 循环按 telemetry_writer_batch_size(默认 200 行)和/或字节预算从两个 deque 批量取行,用一条 INSERT 批量语句写入,最长等待 telemetry_writer_flush_interval_s(默认 0.5s)即落盘半满批次(源码 L702-L751)。
  • 专用引擎 + 微型池:writer 自建一个独立的 AsyncEngine,pool_size 为 1(SQLite)或 2(Postgres),max_overflow=0——从结构上看,遥测流量最多只占用 1~2 条数据库连接,物理上不可能饿死请求池(_create_dedicated_engine)。这正是 2.3 节故障模式消失的根因:writer 关闭时每条遥测 INSERT 都要从请求池借连接,500 并发下池子瞬间见底;writer 开启时生产端零连接占用,只有后台单任务持有小池。
  • 失败重试有上限退避:批量写入失败时整批放回队首重试,指数退避封顶 30s;连续失败达 6 次(_FAILURE_ESCALATION_THRESHOLD)升级为 error 日志,提示持续丢数风险。

4.2 生产端如何把行交给 writer

以 transaction 为例,TransactionService.log_transaction 的逻辑是:

  1. telemetry_writer_enabled=True 且 writer 已 start() → 把整行 model_dump(mode="python") 交给 writer.enqueue_transaction(),入队成功即 return,不获取会话、不占用请求池;
  2. writer 未就绪(如 lifespan 早期)或已停用 → 落到 legacy 路径 async with session_scope() as session: await crud_log_transaction(session, transaction),并只打印一次 "falling back to legacy direct-write path" 警告(用 _legacy_fallback_logged 闩锁防止刷屏)。

vertex_build 侧同构:lfx/graph/utils.py 的 log_vertex_build 在开关命中时经 _try_enqueue_via_telemetry_writer() 入队,否则走 crud_log_vertex_build。harness 直接调用这两个入口,因此它同时覆盖了"入队成功"与"legacy 直写"两条分支。

4.3 磁盘 outbox:进程重启也不丢行

writer 还带一套按 PID 隔离的 SQLite outbox(WAL 模式,<tempdir>/langflow_telemetry_outbox/<pid>/):

  • 优雅关停:teardown() 先等待 writer 排空(受 telemetry_writer_shutdown_drain_s,默认 5s,约束),剩余内存行 append_all 落盘(PRAGMA synchronous=FULL 保证 fsync 后才关连接);
  • 启动恢复:_restore_from_disk 回放本 PID 上次溢出行,_adopt_orphan_outboxes 采纳同主机同 boot 身份(owner.json 中的 hostname + boot_id)的已死进程孤儿目录——防止容器重启后 PID 复用误吞他人数据;跨主机孤儿目录按 owner 文件 mtime 老化(telemetry_writer_orphan_max_age_s,默认 3600s)后剪除;
  • 明确声明的代价:SIGKILL/OOM 这类硬杀会丢失内存中尚未落盘的行——遥测可见性是最终一致而非事务性的,这与这两张表"交互式查询的调试/执行历史"的运维属性匹配。

4.4 容量、字节策略与 retention

所有旋钮集中在 TelemetrySettings 设置组:

配置项缺省值含义
telemetry_writer_enabledTrue是否启用批写入 writer;False 即回退 legacy 直写
telemetry_writer_batch_size200单条批量 INSERT 最大行数
telemetry_writer_flush_interval_s0.5半满批次最长等待秒数
telemetry_writer_max_queue100000每类队列行上限,超限丢最旧行并计数
telemetry_writer_size_strategycountcount/bytes/either 三种阈值策略
telemetry_writer_batch_size_bytes262144字节策略下单批编码字节上限(约一个 TCP 帧)
telemetry_writer_max_queue_bytes209715200(约 200MB)字节策略下队列字节上限,防止单 worker 遥测缓冲主导容器内存
telemetry_writer_shutdown_drain_s5.0关停排空预算
telemetry_writer_cleanup_interval_s60retention sweeper 节奏

retention 不再"每次 INSERT 顺手做",而是由独立的 sweeper 任务摊销执行(_run_retention_pass):按 max_transactions_to_keep(默认 3000)做 per-flow 保留、按 max_vertex_builds_per_vertex(默认 50)做 per-vertex 保留、按 max_vertex_builds_to_keep(默认 3000)做全局保留,且只清理本 flush 周期真正写过的 flow(dirty set 机制),失败时回滚快照以便下轮重试。这也解释了 2.3 节对账为何是"模 retention"的近似相等。

五、把压测用作回归验证的实践建议

  1. 先跑基线(writer 开启,SQLite):200 并发 × 15 秒应得到 error_total=0;记录 flushed_rows 与行增量的偏差,偏差应只来自 retention。
  2. 再跑故障复现(writer 关闭 + 缩池):500 并发 × 15 秒下,legacy 路径应在 tx: / vb: 计数中暴露 OperationalError(SQLite 锁)或 TimeoutError(Postgres 池超时)——这就是 2.3 节环境变量的用途,也是验证"修复确实改变了行为"的对照组。
  3. 换 Postgres 重复上述两步:重点观察 pool_timeout=3 下错误类型是否为池等待超时,而非 SQL 层错误。
  4. 检查 writer 状态行:failed_batches 与 dropped_* 应为 0;若 dropped_tx/dropped_vb 非 0,说明队列溢出丢数,应调大 telemetry_writer_max_queue 或改用 either 字节策略(payload 体积不均的部署推荐,见 设置组注释)。
  5. 注意排空宽限:DRAIN_GRACE_S(缺省 5s)过短会看到 post-drain 时 tx_buffer 仍有残留,这不是丢数——shutdown 路径会把它们 spill 到磁盘 outbox;但行对账 delta 可能因此偏小,必要时调大该变量再测。

最后说明适用前提:该 harness 依赖仓库内 langflow 服务的完整初始化(initialize_services + Alembic 日志文件 + 建表),因此需要在已安装 uv 依赖的项目环境里运行;它验证的是遥测写入子系统的吞吐与隔离性,不覆盖 flow 执行、模型调用等其它路径——这与其文档声明的"无需 OpenAI key 或真实 flow"的定位一致。

【免费下载链接】langflow Langflow is a powerful tool for building and deploying AI-powered agents and workflows. 【免费下载链接】langflow 项目地址: https://gitcode.com/GitHub_Trending/la/langflow

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐