1. 为什么我要从零搭一个轻量 AI 数据平台
去年下半年,我手上同时压着三个跟 AI 沾边的需求:一个要做用户行为数据的特征回填,一个要把业务侧的日志清洗成训练样本,还有一个要给算法同事做一套可复用的数据预处理流水线。三个需求看起来各干各的,但拆开一看,底层干的是同一件事——把散落在各处的原始数据,经过一系列加工,变成 AI 能直接吃进去的格式。
一开始我想得很简单,直接上现成的大数据全家桶不就完了。结果一评估,光是集群资源、组件依赖、运维成本就够我喝一壶的,而且我们团队当时就两三个人,根本没精力去维护一套重型分布式架构。更现实的问题是,业务数据量其实并不大,每天增量也就几十个 G,用重型方案属于拿高射炮打蚊子。
所以我决定自己搭一个轻量级的 AI 数据平台初版。核心目标很明确:用最小的资源开销,把 ETL 流程、算子编排、任务调度这几件事串起来,让数据从原始态到可用态的过程变得可配置、可复用、可观测。这里的关键词是 ETL、算子和架构,后面我会围绕这三个点展开讲。
这篇文章适合谁看?如果你也是小团队里那个"什么都得自己上"的人,或者你正在纠结要不要为了 AI 数据链路引入一套重型框架,那我的这些经验应该能帮你少走点弯路。我不会只讲我做了什么,更会讲我为什么这么选、哪些地方踩了坑、以及这个初版离真正的生产环境还差多远。
先说结论:这个平台初版我大概花了两周多的业余时间搭起来,跑通了从数据接入、算子编排、任务执行到结果落库的完整链路。但它离生产可用还有明显差距,这个差距我在最后一节会掰开揉碎讲清楚。
2. 整体架构设计与技术选型思路
2.1 先想清楚:轻量到底轻在哪
"轻量"这个词很容易被滥用,所以我先给它下个定义。在我这里,轻量意味着三件事:部署简单、依赖少、单机可跑。不需要 Kubernetes,不需要 ZooKeeper,不需要一堆中间件互相注册发现。一个进程能起来,一个配置文件能说清楚所有事,这就是我要的轻量。
那为什么不用微服务架构?我认真考虑过。微服务的好处是解耦、独立部署、独立扩缩容。但在我这个场景下,微服务的代价远大于收益。两三个人的团队,维护五六个服务的注册、配置、链路追踪,光是运维心智负担就压垮了。而且数据平台的任务之间本身就有强依赖关系,硬拆成微服务反而会让调用链变得复杂。
所以我最终选的是单体 + 模块化的架构。一个主进程,内部按职责划分成几个模块:接入层、编排层、执行层、存储层。模块之间通过清晰的接口通信,未来真需要拆分的时候,沿着接口切就行。这是我一贯的做法——先单体跑通,等痛点真的出现了再拆,而不是一开始就假设自己需要分布式。
2.2 核心模块划分
整个平台我拆成了四个核心模块,每个模块的职责边界我卡得很死,避免后面互相污染。
接入模块负责对接各种数据源。可能是数据库、可能是对象存储里的文件、也可能是上游服务推过来的消息。这一层的核心抽象是"数据源连接器",每种数据源实现同一套接口,上层不关心数据具体从哪来。
编排模块是整个平台的大脑。它负责解析任务定义,把任务拆成一个个算子节点,然后按照依赖关系决定执行顺序。这里我用的是有向无环图(DAG)的思路,每个任务就是一张图,节点是算子,边是数据流向。
执行模块负责真正干活。它拿到编排模块给的执行计划,逐个调度算子运行,管理算子的输入输出,处理失败重试。这一层要考虑并发、资源隔离、超时控制这些事。
存储模块负责两件事:一是存任务的元数据(任务定义、执行状态、调度记录),二是存算子处理过程中的中间数据和最终结果。元数据我用的是关系型数据库,中间数据根据体量选择落本地文件系统或者对象存储。
2.3 为什么是这套技术栈
技术选型这块我纠结了挺久,最后定下来的组合是 Python + FastAPI + PostgreSQL + 本地文件系统 + 自研调度器。逐个说下理由。
选 Python 是因为 AI 生态基本都在 Python 这边,算子实现、数据处理库、模型调用,用 Python 最顺手。虽然 Python 在性能上不是最优解,但数据平台的瓶颈通常在 IO 和算子本身的算法复杂度上,语言层面的性能差异不是主要矛盾。
FastAPI 用来做平台的 API 层,主要是因为它异步支持好、开发效率高、自带文档。任务提交、状态查询、算子注册这些接口用它写起来很快。
PostgreSQL 存元数据,这个没什么好纠结的,关系型数据用关系型数据库,稳定可靠,而且它自带的 JSON 字段类型存任务定义这种半结构化数据很方便。
中间数据用本地文件系统,是因为初版阶段数据量不大,而且本地读写延迟低。这里我留了个抽象层,未来换成对象存储只需要改一个实现类。
调度器我是自研的,没用 Airflow 或者 Prefect 这类现成方案。原因是我需要算子级别的细粒度控制,现成调度器的抽象层级偏高,我要在算子之间做数据传递和状态管理,用它们反而要绕很多弯。自研调度器核心逻辑其实不复杂,一个拓扑排序加一个执行队列就够了。
3. 算子体系:整个平台最核心的抽象
3.1 算子到底是什么
算子这个词在不同语境下含义差别很大。在图像处理里,拉普拉斯算子、Halcon 算子指的是具体的卷积核或者图像操作函数。在深度学习里,算子往往指计算图上的一个计算节点。在我这个数据平台里,算子是一个最小的、可复用的数据处理单元。
打个比方,如果把数据处理流程比作做菜,那算子就是一个个具体的操作步骤:洗菜是一个算子,切菜是一个算子,下锅炒是一个算子。每个算子有明确的输入、明确的输出、明确的处理逻辑。你把算子按顺序串起来,就成了一道菜;把算子按不同方式组合,就能做出不同的菜。
这个抽象的好处是极强的复用性。比如"字段重命名"这个算子,在用户行为数据清洗里要用,在日志处理里也要用,我只需要实现一次,两边都能调。再比如"缺失值填充"算子,换个填充策略就是一个新配置,不用重写代码。
3.2 算子的接口设计
算子接口我改了三版才定下来。第一版太简单,只有一个run方法,结果发现没法处理配置、没法上报进度、没法做资源预估。第二版又太复杂,塞了一堆生命周期钩子,写个简单算子要实现的接口太多,开发体验很差。
最终定下来的接口是这样的:每个算子必须实现setup、execute、teardown三个方法,外加一组描述性属性。
class BaseOperator: name: str category: str input_schema: dict output_schema: dict resource_hint: dict def setup(self, config: dict, context: dict): """算子初始化,读取配置,准备资源""" pass def execute(self, inputs: dict) -> dict: """核心处理逻辑,输入输出都是字典""" pass def teardown(self): """清理资源,释放连接""" passsetup负责读配置和准备资源,比如数据库连接、模型加载。execute是核心逻辑,输入是一个字典,键是上游算子的输出名,值是实际数据。teardown负责清理,避免资源泄漏。
input_schema和output_schema这两个属性很关键,它们描述了算子的输入输出结构。编排模块在执行前会校验上下游算子的 schema 是否匹配,不匹配直接报错,避免跑到一半才发现数据对不上。这个设计帮我省了大量调试时间。
resource_hint是给调度器看的,告诉它这个算子大概需要多少内存、是不是 CPU 密集、能不能并发。初版阶段这个 hint 还比较粗糙,但已经能让调度器做一些基本的资源避让了。
3.3 算子的分类与内置算子
我把算子按职责分成了几大类,每类下面实现了一批常用算子。
接入类算子负责从外部读数据,比如ReadFromDatabase、ReadFromCSV、ReadFromObjectStorage。这类算子的特点是输出通常是数据流或者分页迭代器,避免一次性把大表读进内存。
转换类算子是数量最多的一类,负责各种数据加工。字段级的比如RenameField、DropField、CastType、FillMissing;行级的比如FilterRows、SampleRows、Deduplicate;还有聚合类的GroupBy、Join、WindowAggregate。
AI 类算子是跟普通 ETL 平台拉开差距的地方。比如Tokenize做文本分词,Embedding调模型生成向量,FeatureExtract做特征工程,DataAugment做数据增强。这类算子往往依赖外部模型服务,所以我在设计上让它们支持异步调用和批量处理。
输出类算子负责把处理结果写出去,比如WriteToDatabase、WriteToParquet、WriteToFeatureStore。
这里有个设计上的取舍我想特别说一下。我一开始想把所有算子都做成纯函数式的,输入输出完全确定,不依赖任何外部状态。但 AI 类算子天然依赖模型,模型加载是有状态的,硬做成纯函数反而别扭。所以最后我允许算子持有状态,但要求状态必须在setup里初始化、在teardown里释放,保证算子的生命周期是可控的。
3.4 算子编排与 DAG 构建
单个算子没意义,算子串起来才有价值。编排这块我用的是 DAG,任务定义就是一个 JSON,描述有哪些算子节点、节点之间怎么连、每个节点的配置是什么。
{ "task_name": "user_behavior_feature", "nodes": [ {"id": "read", "operator": "ReadFromDatabase", "config": {...}}, {"id": "clean", "operator": "FilterRows", "config": {...}, "depends_on": ["read"]}, {"id": "fill", "operator": "FillMissing", "config": {...}, "depends_on": ["clean"]}, {"id": "embed", "operator": "Embedding", "config": {...}, "depends_on": ["fill"]}, {"id": "write", "operator": "WriteToParquet", "config": {...}, "depends_on": ["embed"]} ] }编排模块拿到这个定义后,先做三件事:校验 DAG 有没有环、校验相邻节点的 schema 是否兼容、计算每个节点的入度和出度。然后跑拓扑排序,得到一个合法的执行顺序。
执行的时候,调度器维护一个就绪队列,入度为 0 的节点先进队列。每执行完一个节点,就把它的下游节点入度减一,减到 0 就进就绪队列。这样天然支持并行——同一层的节点如果没有依赖关系,可以同时跑。
这里有个细节值得说:数据在算子之间怎么传递。我一开始想的是每个算子把输出写到磁盘,下游算子从磁盘读。这样简单,但 IO 开销大。后来改成小数据走内存、大数据落磁盘的混合策略。具体阈值我设的是 100MB,超过就落盘,不超过就直接在内存里传。这个阈值不是拍脑袋定的,是我压测了几轮之后根据内存占用和 IO 延迟的平衡点选的。
4. 实操过程:从零到跑通第一条流水线
4.1 环境准备与项目骨架
环境这块没什么花哨的,Python 3.10 以上,PostgreSQL 14 以上,剩下的都是 pip 能装的库。我强烈建议用虚拟环境,不然依赖冲突能折腾死人。
python -m venv venv source venv/bin/activate pip install fastapi uvicorn sqlalchemy psycopg2-binary pandas pyarrow pydantic项目骨架我按模块划分目录,每个模块一个包,包内再按职责分文件。这个结构看起来普通,但好处是边界清晰,新人接手能快速定位代码。
ai_data_platform/ ├── connectors/ # 数据源连接器 ├── operators/ # 算子实现 │ ├── base.py │ ├── io_ops.py │ ├── transform_ops.py │ └── ai_ops.py ├── orchestrator/ # 编排与调度 │ ├── dag.py │ ├── scheduler.py │ └── executor.py ├── storage/ # 存储层 ├── api/ # FastAPI 接口 └── config/4.2 实现第一个算子
我建议从最简单的算子开始,跑通链路之后再逐步加复杂度。我实现的第一个算子是FilterRows,按条件过滤行。
class FilterRows(BaseOperator): name = "FilterRows" category = "transform" input_schema = {"data": "dataframe"} output_schema = {"data": "dataframe"} resource_hint = {"memory_mb": 256, "cpu_intensive": False} def setup(self, config, context): self.condition = config["condition"] self.engine = config.get("engine", "pandas") def execute(self, inputs): df = inputs["data"] result = df.query(self.condition) return {"data": result} def teardown(self): pass这个算子很简单,但包含了算子该有的所有要素:schema 声明、配置读取、核心逻辑、资源提示。写完这个,我就知道整个算子体系能不能跑通了。
4.3 调度器核心逻辑
调度器是整个平台最烧脑的部分。我实现的是一个基于拓扑排序的并发调度器,核心逻辑大概一百多行。
def schedule(dag, executor, max_workers=4): in_degree = {node.id: len(node.depends_on) for node in dag.nodes} ready = [node for node in dag.nodes if in_degree[node.id] == 0] running = {} results = {} with ThreadPoolExecutor(max_workers=max_workers) as pool: while ready or running: while ready and len(running) < max_workers: node = ready.pop(0) future = pool.submit(executor.run_node, node, results) running[future] = node done, _ = wait(running.keys(), return_when=FIRST_COMPLETED) for future in done: node = running.pop(future) results[node.id] = future.result() for downstream in node.downstream: in_degree[downstream.id] -= 1 if in_degree[downstream.id] == 0: ready.append(downstream) return results这段代码看着简单,但有几个坑我踩过。第一个坑是异常处理,一开始我没处理算子抛异常的情况,结果一个节点失败整个调度器就卡死了。后来改成捕获异常、标记节点失败、把下游节点全部标记为跳过。第二个坑是资源竞争,多个算子同时跑的时候内存爆了,后来加了基于resource_hint的资源感知调度,内存密集的算子限制并发数。
4.4 跑通第一条完整流水线
环境搭好、算子写完、调度器跑通之后,我拼了一条完整的流水线来验证:从数据库读用户行为数据,过滤掉无效记录,填充缺失字段,调 Embedding 算子生成向量,最后写到 Parquet 文件。
第一次跑的时候问题一堆。数据库连接超时、字段类型不匹配、Embedding 服务限流、Parquet 写入 schema 冲突,几乎每个环节都出过问题。但正是这些问题让我把平台的健壮性一点点补起来了。比如连接超时让我加了连接池和重试机制,schema 冲突让我把 schema 校验提前到了执行前。
跑通那一刻其实没什么激动人心的,因为我知道这只是开始。真正难的是让它稳定、让它能处理各种边界情况、让它离生产可用更近一步。
5. 踩过的坑与常见问题排查
5.1 算子设计上的三个典型坑
第一个坑是算子粒度过细。我一开始把每个小操作都做成独立算子,结果一个简单任务要串十几个算子,编排配置长得吓人,调试起来也痛苦。后来我把一些高频组合封装成复合算子,比如CleanAndFill把清洗和填充合成一个,配置量直接减半。经验是:算子的粒度应该以"业务上不可再分的处理单元"为准,而不是以"代码上最小的操作"为准。
第二个坑是算子状态管理混乱。有个算子我图省事,在execute里缓存了模型对象,结果并发执行的时候多个线程共享同一个模型实例,出了诡异的并发问题。后来强制要求所有状态必须在setup里初始化,并且每个执行实例持有独立的状态副本。
第三个坑是 schema 演进。上游算子输出加了个字段,下游算子没适配,跑起来才发现。后来我在 schema 校验里加了兼容性检查,上游新增字段是允许的,但删除或改类型必须显式声明,否则直接拒绝执行。
5.2 调度与并发问题速查
| 问题现象 | 可能原因 | 排查方向 | 解决方案 |
|---|---|---|---|
| 任务卡住不动 | 有环或依赖未满足 | 检查 DAG 是否有环、上游是否失败 | 加环检测、失败传播机制 |
| 内存暴涨 | 算子并发过高或数据未分页 | 看 resource_hint 和实际内存 | 限制并发、强制分页读取 |
| 任务重复执行 | 调度器状态未持久化 | 检查重启后状态恢复逻辑 | 元数据落库、加幂等键 |
| 算子超时 | 单算子处理数据量过大 | 看算子输入数据规模 | 拆分算子、加超时控制 |
| 结果不一致 | 并发写同一资源 | 检查是否有共享写入目标 | 加锁或改为串行写入 |
这张表是我踩坑踩出来的,每一条背后都有真实的事故。特别是"任务重复执行"这条,我有次重启服务,结果之前跑了一半的任务全部重新跑了一遍,把下游数据写重了。后来加了幂等键和状态持久化才解决。
5.3 几个容易被忽略的细节
日志要带任务 ID 和节点 ID。不然任务一多,日志混在一起根本没法排查。我一开始日志就打一行文本,后来改成结构化日志,每条都带task_id、node_id、timestamp,排查效率提升明显。
中间数据要有清理策略。跑得多了磁盘会被中间文件塞满。我加了个定时清理任务,超过 7 天的中间数据自动删除,重要结果才长期保留。
配置要支持环境变量覆盖。本地开发、测试、生产环境的配置肯定不一样,硬编码配置是灾难。我让所有配置都支持环境变量覆盖,部署的时候不用改代码。
提示:算子开发有个反直觉的经验——越是简单的算子越要写测试。因为简单算子复用率高,一个 bug 影响面极大。我那个
FillMissing算子就因为一个边界条件没处理,导致一批数据的填充结果全错,排查了半天。
6. 初版与生产环境的真实差距
6.1 差距一:可靠性与容错
初版能跑通,但离"可靠"还差得远。生产环境要求的是:任务失败能自动重试、节点宕机能故障转移、数据能保证 exactly-once 语义。我这初版只做到了失败标记和手动重跑,自动重试和故障转移都还没有。
exactly-once 这块尤其难。数据平台最怕的就是重复处理或者漏处理。要做到 exactly-once,需要算子支持幂等、需要状态持久化、需要事务性的输出提交。这三样我目前只做了状态持久化,另外两样是下一步的重点。
6.2 差距二:可观测性
初版只有基础日志,生产环境需要的是完整的可观测性:指标(metrics)、日志(logs)、链路追踪(traces)三件套。我现在能回答"任务成功还是失败",但回答不了"每个算子耗时多少、瓶颈在哪、资源利用率如何"。
下一步我打算接入 Prometheus 做指标采集,每个算子执行时上报耗时、内存、CPU 使用情况。链路追踪用 OpenTelemetry,把一次任务执行的全链路串起来。这些做完,排查问题才能从"猜"变成"看"。
6.3 差距三:性能与规模
初版是单机的,数据量一大就顶不住。生产环境要考虑分布式执行、数据分片、算子并行。这块我暂时不打算动,因为当前数据量还没到瓶颈。但我留了扩展点:调度器的执行层是抽象的,未来换成分布式执行器只需要实现同一套接口。
性能上还有个隐性问题是算子本身的效率。我有些算子实现得比较朴素,比如Deduplicate用的是全量加载去重,数据量大了内存扛不住。生产环境需要改成基于哈希分片的流式去重。这类优化得逐个算子过一遍。
6.4 差距四:数据质量与治理
初版对数据质量基本没管,脏数据进来就进来了。生产环境需要数据质量校验、血缘追踪、元数据管理。数据质量校验我打算做成算子,在关键节点插入校验逻辑,不通过就阻断。血缘追踪需要记录每个算子的输入输出来源,这个对排查数据问题极其重要。
6.5 差距五:安全与权限
初版没有任何权限控制,谁都能提交任务、谁都能看所有数据。生产环境需要任务级、数据级的权限控制,需要审计日志。这块涉及的东西比较多,我打算放到后面做,但架构上要预留好扩展点,别到时候改不动。
7. 我对这个初版的一些真实体会
搭这个平台最大的收获不是代码本身,而是对"数据平台"这件事的理解深了一层。以前我觉得数据平台就是把 ETL 工具包一层,真做起来才发现,核心难点从来不是单个功能怎么实现,而是这些功能怎么组合成一个稳定、可观测、可演进的系统。
算子这个抽象是我最满意的设计。它把复杂的数据处理流程拆成了可复用、可测试、可组合的单元,让整个平台的扩展性好了很多。新加一种数据处理逻辑,实现一个算子就行,不用动其他任何地方。
但我也清楚,初版就是初版,它解决的是"从 0 到 1"的问题,不是"从 1 到 100"的问题。生产环境要的那些东西——可靠性、可观测性、性能、治理、安全——每一样都是硬骨头,都得单独花大力气啃。
如果你也在搭类似的东西,我的建议是:先把核心抽象定下来,把最小闭环跑通,别一上来就追求生产级。抽象定错了,后面改起来伤筋动骨;闭环没跑通,你根本不知道真正的痛点在哪。等闭环跑通了,再根据实际暴露的问题逐个补强,这个顺序不能反。
最后分享一个我踩坑踩出来的小技巧:给每个算子写一个最小可运行示例。不用复杂,就是几行代码,构造输入、跑算子、打印输出。这个习惯帮我快速验证算子逻辑,也方便后来人理解算子怎么用。文档可以骗人,能跑起来的示例不会。