数据采集生产系统:从一个采集功能到覆盖多角色、多供应商、多地域的生产体系

已上线

起点是老板一句"这个月要采到 100 小时",围绕这个目标从 Research Box 的采集模块里抽出一条端到端的采集链路做出 MVP;随后被推到 1000 小时、引入供应商和运营,逐步演变成公司规模最大的系统。中间经历过一次由 CTO 主导的整体重构初版,我接手后花大约两周把它从"演示可用"拉到"生产可用";再往后一层,从"能生产"到"能规模化生产"的关键是自动质检流水线,而质检流水线背后又推着我重新想清楚了整套计算基础设施——Flight + Ray 常驻服务模式就是在这个过程里落下来的。这个案例把这两条主线放在一起讲。

在公司的整个技术栈里,这个系统承载的角色比较特殊。它既是业务系统——覆盖平台管理员、供应商管理员、采集员、质检员多角色,跑真实的生产任务、真实的结算;又是数据基础设施——所有采集数据从这里入库、走质检、走后续处理、进入下游标注和训练准备。系统上线之后,公司所有的数据采集都跑在这套体系上,国内和马来西亚两地都有真实生产。

整个项目一年下来我的投入大致是这样一个量级:整体贡献约 70%,代码贡献约 70%,架构设计贡献超过 50%。Demo 阶段主要由 Leader 负责,中间有一段时间下属协助,后续的生产化、规模化和演进阶段我是主要技术负责人。

章一:从采集到生产体系

从采集功能到生产管理系统

处境

Research Box 封存之后,公司不再追求一步到位的完整平台。但用户侧真正持续出现的需求非常明确:数据采集。老板给了第一个具体目标——"这个月要采到 100 小时"。这个目标下面还带着一整套隐含要求:什么时候把项目建好、什么时候把人拉进来、什么时候开始采、每天采多少、怎么保证数据能被后续消费。

当时的现实条件是这样的:真正的采集端还完全没有做;原本 Research Box 的采集能力是围绕机器人接入设计的,不是给"人拿设备去采"的场景准备的;原本的计划是以自研硬件为主采集设备,Infra 同学在设计一套完整的 MQTT 登录认证交互协议,硬件侧负责设备端实现;iPhone 只是一个 backup 选项,不是主力。

桌上的选项

一条路是等硬件——继续按最初的计划推进 MQTT 协议和自研设备的设备端实现,等硬件团队把设备做出来再启动采集。这条路的问题是时间不可控,"这个月 100 小时"这个目标就悬空了。另一条路是不等硬件,先在 Research Box 现有采集模块上做渗透式改造——把原来面向机器人主动接入的接口改造成通用的"数据创建、任务领取、数据上传"三类接口,然后在 iOS 端补齐一整套采集交互,把手机变成实际的采集端。

我的选择

选了第二条。原有 Research Box 的采集模块提供了改造入口:把针对机器人的主动提交接口改造成通用的数据创建 / 领取 / 上传接口;iOS 端从零做了完整的采集能力——登录界面、任务列表、任务详情、任务领取、录制流程、上传流程,做成一整套"用户拿手机就能完整走完一次采集"的交互。整个 MVP 大家还是合作去做的,但因为主要精力都放在自研硬件的 MQTT 协议设计上,Web 前后端和 iOS 这一块的设计和实现基本都是我在做。

设备侧走了一个从"计划走硬件"到"实际走 iPhone"的转向。原本设计的 MQTT 交互协议由 Infra 和硬件团队负责,最终因为硬件没有按期完成,那一套完整协议没有被真正采用;生产目标最后是用 iPhone 完成的。这个转向本身也说明了一件事——最初拆解目标的时候没有把方案绑死在硬件上,采集链路的抽象层是设备无关的,所以硬件掉链子的时候,iPhone 能顶上而不用推翻。

代价

MVP 阶段的功能是收敛的——只有采集员一个角色,只有单一采集端,供应商、运营、多级质检这些都还没有;MQTT 那一套完整的自研协议没有真正落地,前期设计的投入没有直接产出。

结果

100 小时的目标在这个 MVP 上被跑通了。整个链路——登录 → 查看任务 → 领取 → 采集 → 上传——第一次真正端到端串起来。

从"采集功能"到"生产管理系统"的转折

MVP 跑完之后,老板把目标从 100 小时推到 1000 小时。这一步表面上看只是"量的扩大",实际改变的是系统的定位。1000 小时意味着不可能只靠内部人员去采,必须引入供应商;引入供应商就必须引入运营去管理供应商;有了多方参与就必须有任务分发、进度追踪、质量反馈、结算依据。系统从"给采集员用的采集工具"变成了"给平台、运营、供应商三方共用的生产管理系统"。

这个阶段还有一件很关键的事——为了保证数据质量,全员对库里已有的几百个小时数据做了一轮抽检。就是在这次抽检里,团队第一次把"数据必须经历一系列的检查流程才能作为标准入库"这个原则明确下来。这个原则后来变成整个系统里质量体系的基座——所有后续的自动质检、人工质检、批次审核、多级复核都是建立在这条原则之上的。质量标准不是从设计文档里推出来的,是从真实数据里看出来的。

供应商自治作为架构原则

在这个阶段还有一个决策后来被反复证明是关键的:供应商自治。

平台管理供应商,但不管供应商的一线人员。每个供应商被创建之后,自行管理自己的采集员账号和质检员账号,自行分配任务给下属。平台只面对供应商这一层做任务派发、质量反馈和结算,不下沉到具体的采集员。这个设计从一开始就是这样定下来的——不是后期为了扩展性补出来的——原因很直接:如果平台管所有一线人员,那就变成了"平台化",供应商没有自己管理业务的能力,协作方式会非常僵。真的接了几家供应商跑起来之后,这个判断被验证:供应商内部的排班、人员轮换、绩效管理都是他们自己的事,平台把这些抽象成一个统一的"供应商内部管理体系"交给供应商自己用,大家的边界清晰,运营成本也低。

海外业务与系统重构

差不多同一时间,产品和融资侧确认了要和一个海外业务运营方合作。老板安排了一次针对整体系统 (包括 iOS 和 Web 前端) 的样式重设计,并对供应商管理侧提出了一系列新诉求。时间紧、任务重,大家评估之后决定不在原有系统上继续迭代,直接重新设计一套业务系统。

这次重构的收敛范围我在设计讨论的时候明确提过——保留存储能力、保留手机端交互接口,只重构业务系统上层——不动数据底座,不动已经在生产环境里跑起来的采集端。用户管理、项目流程、供应商管理、抽检流程这些抽象成独立业务系统重做,把这一层"关于人和流程的复杂度"从原来的采集系统里剥出来。

批次账单——从平台/供应商信任问题反推出来的机制

这一节是这套生产系统里我个人最想强调的一块。它不是一个功能,是一整套解决"平台和供应商之间数据质量共识"的机制设计,而且它不是拍脑袋想出来的,是从真实的业务问题里反推出来的。

问题一开始就非常具体:供应商侧经常对数据质检结果不满意。典型场景是——供应商说他们采了 8 个小时,平台系统算出来最终有效的只有 4 到 6 个小时,双方就开始扯皮。核心矛盾不是"有效数据到底是多少",而是双方对"数据到底怎么被判定为有效"缺少共同认可的事实基础。之前的对账方式是这样的:我从数据库里导一份 Excel 出来给运营,运营拿着这份 Excel 去找供应商,供应商一头雾水——他们没法从 Excel 里看到自己每一条数据是被哪个环节判定为不合格的、什么原因不合格、通过率的分子分母到底是怎么算的,只能被动接受平台给出的结论。

在和运营讨论的时候,运营先提出了一个业务侧的做法——抽检和打回机制:从供应商的批量数据里抽出一批,做检查,按通过比率向供应商证实结款。这个做法本身是解决对账问题的一个抓手,但我在这基础上进一步想到——能不能把"共同事实"这件事在系统里直接做出来?也就是说,不只是在打回场景里让供应商看到事实,而是整个批次的生命周期里,双方看到的数据都是同一份、状态是同步的。

桌上的选项

一条是继续沿用 Excel 对账 + 人工沟通,只在流程上做规范。这条路的问题是它没有从根本上解决信息不对称——供应商永远是"事后看结果",没法在生产过程中就跟进质量。另一条是把批次这个东西做成系统里的一等公民:供应商把想要提交的一批数据显式打包成一个批次,平台对这个批次做二次审核,在整个审核过程中,批次的状态、明细、每一条数据的审核情况、通过率、不通过原因、最终认可数量对双方完全可见。

我的选择

选了第二条,并且把它做成了完整的机制。核心是一条链:生产事实 → 质量判定 → 结算事实

生产事实这一层是采集员实际交付了什么——多少条数据、每条数据的元信息、是哪个采集员在什么时间用什么设备采的。质量判定这一层是每一条数据在自动质检和人工质检里的具体结论——通过还是不通过,不通过的原因是什么,是哪个质检员判定的。结算事实这一层是把前两层的数据汇总成一个批次的结算依据——通过率、可结算的数据量、可结算的金额。三层数据在系统里是显式关联的,任何一层的结论都能追溯到上一层的具体数据。供应商侧和平台侧看到的批次视图是完全一致的——不是"平台一份、供应商一份、内容差不多",而是同一份数据的同一个视图。

代价

批次审核这个流程本身增加了业务侧的操作步骤——供应商需要显式把数据组织成批次,不能像以前那样"随便交、随便算"。系统里的数据状态机也变复杂了,一条数据从采集到最终结算要经过多个显式状态,每个状态之间的流转都要有对应的审核动作。

结果

这个机制在真实业务里已经跑起来了。大概涉及 8 到 12 家供应商,累计有 9 万多条数据、大约 5 万小时进入了通过这个机制的真实结算。从我作为研发的观察来看,这个机制上线之后没有再出现过之前那种"我采了 8 小时你只算 4 小时"的结算争议。以前运营需要拿着 Excel 反复解释,现在双方直接在系统里看同一份批次,争议的发生场景本身消失了。

回看

这件事让我认识到,业务系统里真正难的很多时候不是"实现某个功能",而是把双方之间的隐性冲突显性化,再用数据把它固定下来。批次账单表面上是个统计功能,底层是一个信任机制——它让平台和供应商共享同一组事实,争议就没有发生的空间了。这个思路后来在系统里的其他协作场景(平台和运营、运营和质检员之间的分工)也反复被用到。

技术架构:供应商自治与批次账单数据流

角色体系与权限分层
┌─────────────────────────────────────┐
│              平台侧                  │
│  平台管理员:创建项目、管理供应商、最终确认 │
│  平台质检员:二次复核                  │
└──────────────┬──────────────────────┘
               │ 项目分发
   ┌──────────┼──────────┐
   ▼          ▼          ▼
┌────────┐ ┌────────┐ ┌────────┐
│供应商 A │ │供应商 B │ │供应商 C │  各自独立管理
│ 管理员  │ │ 管理员  │ │ 管理员  │
│ 采集员× │ │ 采集员× │ │ 采集员× │
│ 质检员× │ │ 质检员× │ │ 质检员× │
└────────┘ └────────┘ └────────┘
  平台不管理供应商内部人员,只管供应商本身

批次账单数据流(信任机制)
  采集员提交数据
       ▼
  供应商质检员一审
       ▼
  供应商管理员复核 → 按批次提交平台
       ▼
  平台质检员二审
       ▼
  ┌──────────────────────────────┐
  │ 批次账单(双方可见同一份)      │
  │  · 数据总量                   │
  │  · 每条审核状态 + 不通过原因    │
  │  · 通过率                     │
  │  · 应结算金额                  │
  │                              │
  │  生产事实 → 质检判定 → 结算事实  │
  └──────────────────────────────┘

接手初版到生产硬化

处境

海外业务需求进来的时候,CTO 决定负责整个业务系统的整体重构,并且在一周之内交付了一版可以演示的系统。重构本身我们讨论过——我参与了方案讨论,并且明确提出过收敛范围的约束:保留原有的存储能力和手机端交互接口,只重做业务上层——但真正的重构工作不是我做的,是 CTO 主导完成的。

一周之后拿到手上的这个初版,可以给海外用户演示,基本的界面和主流程都在。但要真正拿到国内生产环境里替换掉旧系统,还差得比较远。这个判断不是主观评价,是把它放到真实生产流程里一走就能发现的——很多在演示场景下看不到的问题,在真实业务里一下子全暴露了。

桌上的选项

一条路是让 CTO 和相关同学继续迭代新系统,我这边继续维护旧系统,直到新系统达到生产可用状态再切换。这条路的问题是——CTO 的精力不可能长期投在这个系统上,而"生产可用"这个门槛跟"演示可用"之间的差距,从我在旧系统里维护的经验判断,不是几个 Bug 的距离,是一整套隐性 case 需要被覆盖的距离。另一条路是我接手新系统的后续迭代,把它从"能演示"推到"能生产",同时保留旧系统作为国内生产的过渡兜底。

我的选择

选了第二条。当时我一边维护旧系统,一边开始接手新系统的问题清单。第一步的目标非常明确:让新系统能在国内替换掉旧的生产系统。第一步跑下来大概花了两周,这两周里修的都不是简单的代码 Bug,而是一整套业务流程串联层面的问题。举几个具体的:

一个是任务配置不生效。运营在后台配置一个采集任务时能填的字段很多——采集类型、数据要求、质检规则、分发规则、结算规则——但配完之后到采集员那一侧,有些字段是取不到的。查下去发现是配置的存储和读取用了不同的数据路径,后台写一处、前端读另一处,单看代码都对,合起来就是不生效。这个类型的问题在初版里不是一个,是一批。

一个是质检规则不完整。原始设计只考虑了主流程——一条数据被创建、被采集、被质检、被通过 / 被打回——但没考虑任务状态的回滚场景:一条数据被打回之后,采集员重新采集了一份,系统里应该怎么处理原来那份的状态?一个批次已经进入审核之后,如果发现里面某条数据需要撤回,批次应该怎么变?这些边缘 case 在演示里不会走到,在真实生产里每天都会遇到。

一个是自动质检和人工质检的流程串联断裂。设计文档里是"数据 → 自动质检 → 人工质检 → 通过入库",但实际实现里,自动质检的结果没有真正被人工质检环节感知,质检员会看到一堆自动质检本应该已经过滤掉的问题;或者反过来,自动质检还没跑完,人工质检就把数据放行了。这一层不是加个字段能解决的,是要把自动 / 人工两条流水线的状态机对齐,把"什么时候能进入人工质检"这个前置条件明确写进业务规则里。

这两周做的事情表面上是"修 Bug",实际上是把一个还没有经过真实业务打磨的系统,一次次拿到真实场景里去撞,撞出来的每一个问题都对应一条需要被补齐的业务规则或者数据模型。

代价

这段时间里我同时要维护两套系统,精力上比较紧;而且这次接手之后,我实际上继承了新系统后续所有的迭代责任——它从一个"临时接过来修一下"变成了长期的技术 Ownership。

结果

两周之后,新系统具备了替换国内生产系统的能力,之后就是完整替换加持续迭代。国内和海外都跑在这一套上,一直演进到今天成为公司规模最大的系统。

回看

这段经历留下的最深的一个判断是:从"能演示"到"能生产"这段路,靠的不是聪明,是系统性地覆盖 case。一个新系统在演示场景下的完备度,和它在真实生产环境里的完备度之间,差的是所有那些"演示的时候用户不会点进去、真实业务里每天都要走"的边缘路径。这段路没有捷径,只能一个 case 一个 case 把它撞出来、修好、写进业务规则里。

关于 CTO 主导那次重构初版这件事,我想说的是——这两件事(重构初版由 CTO 完成、后续生产化由我完成)在时间线上是两个不同的阶段,贡献边界是清楚的;我在描述这个项目时不会把重构本身写成自己的工作。我的核心贡献发生在把一个可演示的东西持续变成真正能生产、能规模化运营、能结算的系统——这个部分从材料事实上讲是我做的。

章二:计算能力升级

自动质检流水线:从"能生产"到"规模化生产"

新系统替换掉国内生产系统之后,下一个关键节点是自动质检流水线。规模化生产要成立,第一个条件是流程能跑,第二个条件是质检不能靠人堆——如果每一条数据都要靠人工从头看到尾,规模一大质检就成为瓶颈,通过率、生产节奏、结算周期全都会受影响。

自动质检流水线要做的事就是把"明显有问题的数据"在进入人工质检之前先过滤掉,让人工质检只花时间在真正需要人判断的数据上。原本的业务系统设计里是没有自动检查这一环的,只有人工质检。我在人工质检前面加了一个自动质检环节,并把它做成一条硬规则:自动质检必须全部通过,才可以进入人工质检

自动质检当前实现了 5 类检测能力:

第一,YOLO 手部检测。采集数据大部分是采集员的第一视角视频,手部动作是核心内容;如果整个视频里人手完全没有出现,那这条数据几乎肯定是问题数据(比如手机被摆在那儿录空镜)。用 YOLO 的现成模型做快速识别,遍历视频帧,判断手部出现的频率和覆盖比例。

第二,黑帧检测。用逐帧的哈希值或颜色值计算判断一帧是不是"完全黑"——纯黑帧、极暗帧都能通过这个方式识别出来。整片视频里如果有大量黑帧,那可能是设备遮挡、镜头盖没拿掉、或者录制中断。

第三,静态视频检测。这个针对的是"手机放在那儿录一段假数据"这类情况。判断方式是计算相邻帧或者多帧之间的差值,如果整段视频里帧间差值一直很小、几乎没有画面变化,就判定为静态。

第四,YOLO 人脸检测。这一条我要诚实说——准确率是这五项里最低的一项。用 YOLO 或者 Face 模型做人脸识别本身有一定误检和漏检率,尤其是侧脸、遮挡、模糊、小尺寸人脸这些场景。所以在实际使用中,人脸检测的阈值设得比较保守——宁可让一些边缘 case 触发人脸检测告警进入人工复核,也不放过任何一个可能的人脸暴露风险。因为一旦人脸真的漏检暴露出去,合规问题的代价比多花一点人力做复核大得多。

第五,过曝 / 抖动检测。这一类是基础的画质判断,判断整段视频是不是曝光严重不足或过度、是否有大量抖动导致数据不可用。

这五项检测能力里,算法本身都用的是现成的模型和现成的计算方法——YOLO 是开源模型,黑帧和静态检测是标准的视频帧处理算法,人脸检测用现成的 Face / YOLO 模型。我做的不是算法研发,而是把这些能力工程化成一条可自动执行、可分发、可扩容、可观测的质检流水线,并把它嵌入到业务系统的数据流转里,让它跟人工质检形成明确的前置关系。

流水线跑起来之后,系统的规模化生产能力是站住了。峰值大概是每天 3000 采集小时的数据吞吐——这是我在真实生产里见过的最大值,而不是系统的容量上限,因为整个系统的吞吐量在跑到 3000 小时/天的时候还没有触碰到瓶颈。黑屏视频、手部完全缺失的视频、被静态摆拍的视频都能被有效过滤;人脸检测因为准确率的限制,主要作为一个"进入人工复核"的触发器,不作为最终判定。

回看

这一段让我意识到,自动质检的核心价值不是"完全替代人工",而是把人的注意力集中到真正需要判断的数据上。如果所有数据都要人从头看到尾,那规模化生产是不可能的;如果自动质检可以拦住 80% 明显有问题的数据,人工质检就有余量去仔细判断剩下 20% 的边界情况。我在这一层的判断是——用什么模型、准确率多少,不是核心指标;核心指标是这条流水线能不能把整体的生产成本压下来、把人的注意力放到该放的地方。人脸检测那一项虽然准确率有限,但只要阈值设置得当,它就在流水线里发挥了正确的作用。

Flight + Ray 常驻模式

自动质检流水线要跑起来,底下需要一套计算基础设施——它不是普通的业务代码,是要处理大量视频数据、要跑各种模型、要做批量转码和分析的计算密集型任务。我之前在另一家公司做过好几套自动化流水线,主要方案是 Argo Workflows 和 Airflow;这一次的场景对底层计算模式的要求不太一样,最后落到了 Flight + Ray 上——但这个方向不是我提出的,是 Infra 团队提出来的,我最初对它是抵触的。这一节讲的就是从抵触到最终形成一套常驻计算服务模式的完整过程。

处境

之前用 Argo / Airflow 这类传统 Workflow 引擎,基本模型是这样一条链:代码或 YAML 定义 Workflow → 用户请求触发 → 创建 Workflow 实例 → 创建 Pod → 执行任务 → Pod 销毁。每一次新的任务请求都要走一遍这个完整生命周期——创建、执行、销毁、回收。

这个模式在低频、复杂编排、一次性任务场景下没什么问题,但在高频、大规模数据处理场景里逐渐暴露出问题。我之前在另一家公司用 Argo 的经验里,每天要处理约 10 万级 Case,底层资源大概是 5000 到 15000 CPU 的规模;跑久了之后,真正的瓶颈往往不是计算本身,而是 Kubernetes 层的任务生命周期管理成本:每一个任务都要拉一次镜像;每一次 Pod 高频启停都要走一遍 Kubernetes 的调度、挂载、初始化流程;大量 Pod 的创建和销毁会给 Kubernetes 的控制面持续施压;时间长了之后偶尔会出现 Pod 无法被及时删除,残留 Pod 占用资源不释放,后续任务要等资源释放再启动。这些问题单个看都不大,规模化之后就会持续影响整体的资源吞吐。

到了这一家公司,自动质检的场景需求本质上也是高频、大批量、重复性的数据处理——H.264 → MP4 Remux、视频解码、逐帧的黑屏 / 手部 / 静态 / 人脸 / 过曝分析——这类任务如果还沿用 Argo 那种"每个任务起一遍 Pod"的模式,前面那些问题都会重演。

桌上的选项

一条路是继续用熟悉的 Argo / Airflow,把之前踩过的坑做一些工程优化——镜像预热、Pod 池、调度参数调优——这条路是渐进的,但没有本质改变模式,规模再大一档还是会撞到同样的天花板。另一条路是走 Infra 提出的 Flight + Ray——Flight 作为上层的 Workflow 定义和服务入口,Ray 作为下层的分布式计算和弹性 Scaling 引擎——但这个方向我最初是不认可的:Flight 的 API 使用体验并不好,Workflow 层本身有一定限制,Workflow 不适合灵活承载多种不同任务,每一次创建 Workflow 的复杂度也不低。从"选一个更好用的工具"的角度看,Flight + Ray 不比 Argo 更吸引人。

我的选择

真正让我改变判断的是 Infra 同学提供了一个最小 Demo 之后的一次对话。他给了一个关键提示——把 Workflow 的启动和任务的调度顺序做一次交换。原本的用法是"请求来了 → 创建 Workflow → 执行",Workflow 是随任务生命周期起落的;换一种用法是"提前部署 Workflow → Workflow 常驻 → 监听消息队列 → 消息到达时分发到 Ray 执行"——Workflow 本身变成一个长期存在的计算服务入口,任务只是往这个服务里灌进来的输入。

这个视角转换一下子把我之前对 Flight + Ray 的所有抵触点都消解掉了。Flight 的 API 是不是好用不重要了——因为你只需要"部署"它一次,不需要每个任务都跟它交互;Workflow 灵活性不够也不重要了——因为你的 Workflow 是"我是一个视频质检服务",不是"我是一次具体的质检任务";Workflow 创建复杂度高也不是问题了——因为它只创建一次,之后长期存活。这个方向真正要设计的不是"怎么用好 Flight",而是"怎么把 Workflow 作为一个常驻计算服务来组织"

具体到落地,Infra 团队负责的部分是提供 Flight + Ray 的基础组件、Demo 和底层集成能力;Flight 和 Ray 之间的连接是它们本身在开源上就自带的能力,不需要额外做。我独立完成的部分包括:把 Workflow 的运行模式从"按需创建"改造成"常驻服务";引入 Kafka 作为任务触发的消息通道;实现任务分发逻辑,把打进来的任务按类型和参数分发到不同的处理路径;实现任务的可扩展执行机制,借助 Ray 做实际的并发和弹性 Scaling;整体的服务化封装、错误处理、状态管理、观测手段。做完之后拿几个真实的质检 Case 去验证,验证通过之后就正式上线,把这套模式作为采集系统里数据质检和数据处理服务的底层计算依赖。

代价

这个模式不是免费的,几个明确的代价:第一,固定资源占用——常驻的 Flight Pod 本身持续占用计算资源,即使没有任务进来,Workflow 也是活着的。传统模式下"没任务就没 Pod",现在变成了"没任务也有 Pod"。第二,开发复杂度上升——开发人员需要明确区分:什么代码应该在 Flight 这一层执行(任务调度、流程控制、状态管理),什么代码应该交给 Ray 执行(实际的计算密集型函数)。这个边界如果设计不好,代码复杂度会明显增加。第三,学习成本——团队需要理解 Workflow 层和分布式计算层之间的边界,以及 Ray 的编程模型;这个门槛比直接写 Argo YAML 高。

结果

这套模式实际用在采集系统的数据质检、数据预处理、数据后处理链路上,一直跑到今天。从对比的直观感受来讲,Argo 那边遇到的所有和 Pod 生命周期相关的问题——镜像下载等待、Pod 创建销毁的开销、Kubernetes 控制面的压力、残留 Pod 的资源占用——在这套模式下都没有再遇到:Pod 一旦起来就长期在跑,资源始终连续被有效使用,不存在"前一个任务结束 → Pod 销毁 → 后一个任务再拉起 Pod"这段等待时间;Kubernetes 层也不会因为高频 Pod 起落而承担额外压力。我没有做严格的性能 benchmark(两个环境不同,直接对比不公平),但工程上的差异是明确的。

回看

这段经历里我真正想强调的不是"Flight + Ray 是更好的选择",而是——组件本身不是关键,组件的运行方式才是关键。同样的 Flight + Ray,按"每次请求创建 Workflow"的用法,只是一个稍微不同的 Argo;按"Workflow 常驻 + 消息驱动"的用法,才是另一种计算范式。这个判断是我自己在实践里撞出来的,不是从技术文档里读出来的。

而且我最后没有把这套模式当成万能钥匙。训练系统那一侧,数据处理的场景和采集侧完全不一样——训练侧的数据处理典型是"一份数据集 → 大规模归一化 / 结构转换 → 完成",是一次性的、批量的、单数据集的处理,不是高频的流式请求;这种场景用 Flight + Ray 的常驻服务 + Ray Scaling 模式反而是浪费,因为流式弹性能力用不上,常驻服务的固定开销白白付出。所以我给训练系统那一侧提的建议是用相同的 Flight 基建、但换一种任务模板——启动一个 Pod,在这个 Pod 的所有资源之内尽可能地完成整个数据集的处理,不使用 Ray 的流式能力。这个建议后来也落地成了训练系统里的实际任务模板。

这个"按 Workload 特征选择计算范式"的判断,我认为是这段实践里最有价值的一层认知——它不是关于某个工具的,是关于在有多种可选计算模式时,基于工作负载的特征做出合适的选择。低频、复杂编排、一次性任务 → Argo / Airflow;高频、大批量、重复性的流式处理 → 常驻 Flight + Ray;单数据集、批量的大规模处理 → 单 Pod 吃满资源。这个判断在之后的系统设计里反复被用到。

技术架构:计算范式与质检流水线

计算范式选型(按负载特征决定)

高频流式任务(采集质检)          批量单次任务(训练数据处理)
┌──────────────────────┐    ┌──────────────────────┐
│ Flight + Ray 常驻模式  │    │ 单 Pod 拉满资源模式    │
│                      │    │                      │
│ 预部署 Workflow       │    │ 请求 → 创建 Pod      │
│      ↓               │    │      ↓               │
│ 常驻 Pod(不销毁)    │    │ 执行 → Pod 销毁      │
│      ↓               │    │                      │
│ 监听 Kafka 消息       │    │ 适合:大数据集        │
│      ↓               │    │ 单次处理,资源密集     │
│ Ray 任务调度 + 弹性伸缩│    └──────────────────────┘
│      ↓               │
│ 执行完继续监听        │
│                      │
│ 适合:高频连续        │
│ 避免 Pod 冷启动       │
└──────────────────────┘

自动质检流水线(串联在 Flight+Ray 上)
  数据提交
     ▼
  ┌─ 手部检测(YOLO)─── pass?
  ├─ 黑帧检测(帧哈希)── pass?
  ├─ 静态检测(帧间差异)─ pass?
  ├─ 人脸检测(YOLO)─── pass?  ← 准确率有限,保守阈值
  └─ 过曝/抖动检测 ──── pass?
     ▼ 全部通过
  进入人工质检队列

多地域、当前技术债与回看

多地域

国内和马来西亚两地都在跑真实生产,同一套代码基础上通过配置区分地域差异。这个能力从系统设计阶段就考虑到了——不是后期"再部署一套"的补丁,而是把地域差异抽象成配置项。可配置的部分包括:两地的质检规则不同、马来侧需要对数据生产后打特定的数据标签而国内不需要、两地使用的采集设备不同(国内有自研新设备,马来目前只用 iPhone)、检测条件也不完全相同。这些差异都通过运营端的配置化管理,一套代码在两个环境里根据配置正常运行。

跨地域的核心不是把服务部署到两个 Region,是同一套业务系统在两种不同的运营配置下都能稳定工作

当前正在处理的技术债

这个系统一路做大,现在已经演变成大约 20 个微服务的规模。服务多本身不是问题,真正的问题是——因为长期用 Web Coding 快速迭代,服务之间的边界没有被严格定义。一个业务需求进来,经常需要同时改动 4 到 5 个服务的代码;每个研发要维护这个系统,需要对全局有比较高的理解;新人接手成本高。这个是我最近在花精力解决的技术债。

治理思路是:重新划清每个服务的边界,用 gRPC 定义清晰的服务能力和交互协议,同时把 gRPC 和 HTTP 的适用场景分开——gRPC 主要用于服务间的能力调用,HTTP 主要用于面向前端的业务接口。这样一次业务需求就能拆到具体的服务上,每个研发只需要理解他负责的那个模块,就能独立维护。目标不是"微服务做得更漂亮",而是让系统的迭代复杂度不再随规模无限增长,让任意一个研发不需要理解整个系统也能承接工作

这件事我现在还在推进,不是已完成的事实,应当作为"当前正在处理的问题"来讲,而不是当作项目成果。

回看这整个项目

一年下来这个项目的完整轨迹是这样一条曲线——采集功能 → 生产管理系统 → 规模化生产系统 → 需要重新治理复杂度的成熟系统。每个阶段面对的问题都不一样:第一阶段是"怎么把最小的采集链路端到端跑通";第二阶段是"怎么把多角色的生产协作在系统里体系化";第三阶段是"怎么让规模上来之后系统不成为瓶颈";第四阶段是"怎么在系统已经很大之后重新拿回可维护性"。

我在这个项目里的价值主要发生在几个具体判断和实际落地上:MVP 阶段没有把方案绑死在硬件上,让 iPhone 能在硬件掉链子的时候顶上;供应商自治作为架构原则从一开始就明确,让后来接多家供应商没有变形;从运营的抽检提议进一步抽象出批次账单的信任机制,让数据质量的争议在系统层面消失;把 CTO 交付的重构初版从"能演示"推到"能生产",这段没有捷径,只能一个 case 一个 case 撞;自动质检流水线让人的注意力可以集中到该判断的数据上;Flight + Ray 从抵触到理解到独立实现常驻服务模式,并把认知进一步迁移到训练系统。

有些贡献边界我想再明确一次——最初的 100 小时目标是老板提的,不是我提的;自研硬件的 MQTT 协议是 Infra 和硬件团队做的,我没参与实现;整体业务系统重构初版是 CTO 一周内完成的,我接手的是之后的生产化和演进;Flight + Ray 的方向是 Infra 提出的,组件本身是开源能力,我的核心贡献是常驻 Workflow + Kafka + 任务分发这一层的重新设计和独立实现;YOLO、Face 这些模型是现成的,我做的是把它们工程化成质检流水线。这些边界我在写材料的时候会保留,不会把这些工作全部归到自己名下。

这个项目最能证明的东西不是某个具体的技术亮点,而是在一个从零起步、需求不断加码、约束不断变化的场景里,能够持续把系统推着往前走,同时保持每一步的判断是站得住的。规模从 100 小时到 3000 小时/天;组织从只有采集员到覆盖平台、运营、供应商、多层质检员;地域从国内到国内 + 海外;结算从 Excel 对账到 9 万 + 条数据在系统里跑真实的批次结算——这些数字背后不是某一次突破,是一年时间里一个决策接一个决策堆出来的。