ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Langfuse 摄取事件重放(S3 Replay)运维指南:基于 Athena 访问日志与 BullMQ 队列的失败补偿

Langfuse 摄取事件重放(S3 Replay)运维指南:基于 Athena 访问日志与 BullMQ 队列的失败补偿 Langfuse 摄取事件重放S3 Replay运维指南基于 Athena 访问日志与 BullMQ 队列的失败补偿【免费下载链接】langfuse Open source AI engineering platform: LLM evals, observability, metrics, prompt management, playground, datasets. Integrates with OpenTelemetry, LangChain, OpenAI SDK, LiteLLM, and more. YC W23项目地址: https://gitcode.com/GitHub_Trending/la/langfuse当 Langfuse 自身的 Worker 处理流程或 ClickHouse 写入出现故障时已写入 S3 的原始摄取事件文件可能未被正常消费。本文基于仓库中 worker/src/scripts/replayIngestionEvents/README.md 及对应实现 s3-ingestion-event-replay.ts完整讲解定位失败事件 → 生成事件清单 → 本地连接 Redis → 将事件重新注入摄取队列的端到端补偿方案并深入源码说明事件 key 解析、队列分片与三阶段数据流水线的工作原理。背景Langfuse 的 S3 事件上传与消费链路在 Langfuse 的摄取链路中事件并非只存在于内存队列摄取处理器会先将事件体以 JSON 文件形式批量上传至 S3再向 BullMQ 摄取队列投递任务。这一点可以在 packages/shared/src/server/ingestion/processEventBatch.ts 中看到脚本通过getS3StorageServiceClient(env.LANGFUSE_S3_EVENT_UPLOAD_BUCKET).uploadJson(bucketPath, data)将事件按eventBodyId分组后写入 S3packages/shared/src/server/ingestion/processEventBatch.ts#L279-L316随后任务才进入队列等待 Worker 处理。由于 S3 上的事件文件与队列任务存在双写关系当 ClickHouse 处理失败或队列消费异常时S3 中仍然保留着原始事件文件的完整副本。这正是重放Replay能够成立的根基只要能从 S3 访问日志中还原出哪些文件被写入过就可以把这些 key 重新组织成队列任务让摄取流程重跑一遍从而实现不依赖原始 SDK 调用方的数据补偿。仓库在 worker/src/scripts/replayIngestionEvents/s3-ingestion-event-replay.ts 中提供了完整的重放脚本v1 方案本文即围绕该脚本展开。第一步识别需要重放的事件Athena 查询使用 S3 访问日志最可靠的做法是开启 S3 服务器访问日志Server Access Logging然后用 Athena 直接查询这些日志。S3 访问日志会记录每次REST.PUT.OBJECT写入的 bucket 与 key正好覆盖摄取处理器上传事件文件的全部写入动作。若尚未开启访问日志需要先在 AWS 控制台为 Langfuse 事件 bucket 启用 Server Access Logging并准备一个用于存放日志的目标 bucket具体配置步骤可参考 AWS 官方文档Enabling Amazon S3 server access logging与Querying access logs for requests using Amazon Athena。在 Athena 中执行以下查询即可生成一份候选事件清单select operation, key from mybucket_logs where operation REST.PUT.OBJECT AND parse_datetime(requestdatetime,dd/MMM/yyyy:HH:mm:ss Z) BETWEEN parse_datetime(2025-07-09:00:30:00,yyyy-MM-dd:HH:mm:ss) AND parse_datetime(2025-07-09:07:45:00,yyyy-MM-dd:HH:mm:ss) limit 50要点说明operation REST.PUT.OBJECT只筛选对象写入操作即摄取事件文件的落盘动作parse_datetime(requestdatetime, dd/MMM/yyyy:HH:mm:ss Z)S3 访问日志中的时间字段是dd/MMM/yyyy:HH:mm:ss Z格式例如09/Jul/2025:00:30:00 0000需要用该格式解析并与事故窗口对比时间范围BETWEEN ... AND ...应覆盖故障发生到恢复的完整窗口确保不遗漏任何失败期间写入的事件可根据需要把limit 50去掉或调大让查询返回窗口内的全部记录。手动构造 CSV兜底方案如果环境中没有可用的 S3 访问日志也可以根据故障窗口内实际写入的事件自行整理一份 CSV。无论来源如何重放脚本只认以下两列格式operation,key REST.PUT.OBJECT,projectId/type/eventBodyId/eventId.json ...其中operation固定为REST.PUT.OBJECTkey为事件文件在 S3 bucket 中的对象键相对LANGFUSE_S3_EVENT_UPLOAD_PREFIX前缀格式为projectId/type/eventBodyId/eventId.json。生成好 CSV 后必须将其放置在仓库的worker/events.csv即 worker/ 目录下因为脚本默认以相对路径events.csv读取输入文件。第二步准备本地环境变量连接 Redis 等基础设施重放脚本需要直接连接 Redis写入 BullMQ 队列同时启动脚本时还需要完整的 Langfuse 环境变量。在仓库根目录创建一个.env文件内容如下# Relevant REDIS_CONNECTION_STRINGredis://:myredissecret127.0.0.1:6379 LANGFUSE_S3_EVENT_UPLOAD_BUCKETbucket-name # Necessary for parsing the file and starting the script CLICKHOUSE_URLhttp://localhost:8123 CLICKHOUSE_USERclickhouse CLICKHOUSE_PASSWORDclickhouse DATABASE_URLpostgresql://postgres:postgreslocalhost:5432/postgres各变量的作用变量说明REDIS_CONNECTION_STRING目标 Langfuse 实例使用的 Redis 连接串。脚本通过它创建 BullMQ 队列实例并把重放任务写入 Redis必须指向生产环境实际使用的 Redis含密码与地址LANGFUSE_S3_EVENT_UPLOAD_BUCKET事件文件所在的 S3 bucket 名。脚本解析 S3 key、构造bucketPrefix时依赖它见下文事件 key 解析一节同时 env 解析器要求该值存在见 packages/shared/src/env.tsCLICKHOUSE_URL/CLICKHOUSE_USER/CLICKHOUSE_PASSWORD脚本由dotenv -e ../.env -- tsx ...启动加载langfuse/shared时会触发 zod 环境变量校验CLICKHOUSE_URL在 packages/shared/src/env.ts 中为必填项因此需要提供可解析的占位值DATABASE_URL同上属于共享库环境校验所需的 PostgreSQL 连接配置若你的环境启用了 Redis ClusterREDIS_CLUSTER_ENABLEDtrue则还需关注 packages/shared/src/env.ts 中的集群开关以及LANGFUSE_INGESTION_SECONDARY_QUEUE_SHARD_COUNTpackages/shared/src/env.ts它们决定队列分片数量详见队列与分片机制一节。第三步执行重放脚本在仓库根目录运行pnpm run --filterworker refill-ingestion-events该命令在 worker/package.json 中定义等价于dotenv -e ../.env -- tsx src/scripts/replayIngestionEvents/s3-ingestion-event-replay.ts即通过dotenv加载仓库根目录.env再用tsx直接运行 TypeScript 脚本无需预先编译。脚本执行过程中会实时打印进度统计每 10 秒输出一次已处理行数与吞吐率并在每个阶段结束时汇总指标包括CSV 总行数、过滤后行数、JSONL 转换对象数、OTEL 对象数、入队事件数、处理耗时与平均速度等便于运维人员确认补偿进度。故障处理大文件拆分与重跑当events.csv过大导致脚本报错如字符串长度超限时README 给出的标准做法是用 Linuxsplit命令把文件按行数均分为多份split -l $(($(wc -l events.csv) / 4)) events.csv part_要点wc -l events.csv先统计总行数除以 4 得到每个分片的目标行数split -l按该行数切分切分后的文件名为part_aa、part_ab等每个分片都需要保留并更新 CSV 表头即operation,key这一行否则脚本找不到operation/key列会直接报错退出建议控制每个events.csv的总大小在 150MB 左右这是仓库文档给出的经验阈值可规避超大文件带来的内存与解析问题分片处理完毕后依次重命名并逐个运行重放命令即可。深入源码重放脚本的三阶段流水线s3-ingestion-event-replay.ts的核心逻辑是一个过滤 → 转换 → 入队的三阶段流水线main()见 s3-ingestion-event-replay.ts阶段一过滤 CSVfilterCsvFile使用csv-parse以流式方式读取events.csvcolumns: false保持数组行skip_empty_lines: true跳过空行从表头中定位operation与key两列索引任一列缺失即中止仅保留operation REST.PUT.OBJECT的行写入events_filtered.csv每 10 秒打印一次处理进度结束时输出过滤统计与文件体积缩减率。阶段二转换 JSONLconvertCsvToJsonl将过滤后的 CSV 按 S3 key 的形态转换为两类 JSONL标准事件 →events_filtered.jsonlOTEL 事件 →otel_events_filtered.jsonl。转换的核心是 packages/shared/src/server/ingestion/eventBucketPath.ts 中的parseEventKey它通过两个正则区分 key 形态eventBucketPath.ts标准格式: projectId/entityType/eventBodyId/eventId.json OTEL格式: otel/projectId/yyyy/mm/dd/hh/mm/eventId.json标准 key会进一步解析出projectId、entityType、eventBodyId、eventId并生成重放任务载荷type取${entityType}-createbucketPrefix通过rawEventBucketPrefix按 S3 中的字面段原样重建eventBucketPath.ts。这里刻意不做二次清洗不调用safeBlobKeySegment因为解析出的eventBodyId已经是 S3 侧规范字符串任何改写都会导致重放时找不到原始文件——这是源码注释中明确强调的边界eventBucketPath.tsOTEL key只提取projectId以fileKey携带完整 key 路径并标记sdkName/sdkVersion为UNKNOWN_INGESTION_SDK_VALUE无法匹配任何格式的行会被跳过并打印告警日志。转换后的两类载荷结构对应脚本中的JsonOutputItem与OTelJsonOutputItem见 s3-ingestion-event-replay.ts——标准事件写入IngestionSecondaryQueueOTEL 事件写入OtelIngestionQueue。阶段三批量入队ingestEventsToQueue/ingestEventsToOtelQueue以readline逐行读取 JSONL按BATCH_SIZE 1000批量入队s3-ingestion-event-replay.ts每个事件通过SecondaryIngestionQueue.getInstance({ shardingKey })或OtelIngestionQueue.getInstance({ shardingKey })选择目标分片队列shardingKey分别为projectId-eventBodyId与projectId-fileKey同一批次内按队列聚合后调用queue.addBulk(jobs)批量投递s3-ingestion-event-replay.ts每个任务附带randomUUID()生成的id与当前时间戳入队完成后关闭所有队列连接输出汇总统计。队列与分片机制重放任务如何回到摄取链路重放任务最终进入的两个队列在 packages/shared/src/server/queues.ts 中定义IngestionSecondaryQueuesecondary-ingestion-queue标准摄取事件的次级队列用于隔离高优先级/高吞吐项目OtelIngestionQueueotel-ingestion-queueOTEL 摄取队列。两个队列类都实现了基于 shard 的实例管理SecondaryIngestionQueue.getShardIndexFromShardName解析分片索引在REDIS_CLUSTER_ENABLED true且提供shardingKey时通过getShardIndex哈希分配分片否则回退到分片 0packages/shared/src/server/redis/ingestionQueue.ts。队列名形如secondary-ingestion-queue、secondary-ingestion-queue-1……分片数量由LANGFUSE_INGESTION_SECONDARY_QUEUE_SHARD_COUNT/LANGFUSE_OTEL_INGESTION_SECONDARY_QUEUE_SHARD_COUNT控制packages/shared/src/env.ts。此外队列默认作业选项体现了重放任务的容错设计removeOnComplete: true完成后自动清理、removeOnFail: 100_000保留失败任务便于排查、attempts: 5与backoff: { type: exponential, delay: 5000 }指数退避重试。这意味着即便重放时 ClickHouse 仍不稳定任务也会按退避策略自动重试而非一次性失败丢失。补充更轻量的 v2 重放方案若你的环境不便从本机直连生产 Redis/ClickHouse/PostgreSQL仓库还提供了 v2 重放方案 worker/src/scripts/replayIngestionEventsV2/README.md它通过POST /api/admin/ingestion-replay管理端点提交 S3 key 列表只需LANGFUSE_HOST、ADMIN_API_KEY和 Athena 导出的events.csv无需克隆仓库与直连基础设施并内置了断点续传--resume、限速--rate-limit与批次重试等能力。在需要长期维护的运维场景中v2 是 v1 的推荐替代方案而本文所述的 v1 脚本则适合在已有仓库克隆与全套环境变量的前提下直接以pnpm命令完成补偿。小结Langfuse 的 S3 事件重放机制本质上是利用事件文件与队列任务双写的架构冗余完成数据自愈S3 访问日志提供事件清单Athena 提供查询入口refill-ingestion-events脚本完成从 CSV 到 BullMQ 队列的链路重建。掌握这套流程后运维人员可以在不打扰原始 SDK 调用方的前提下安全、可控地补偿故障窗口内的全部摄取事件——这既是 Langfuse 摄取链路可靠性的最后一道保险也是理解其 S3 事件存储与分片队列设计的最佳切入点。【免费下载链接】langfuse Open source AI engineering platform: LLM evals, observability, metrics, prompt management, playground, datasets. Integrates with OpenTelemetry, LangChain, OpenAI SDK, LiteLLM, and more. YC W23项目地址: https://gitcode.com/GitHub_Trending/la/langfuse创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进