news 2026/10/7 17:52:43

从零搭建轻量AI数据平台:ETL算子编排与DAG调度实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
从零搭建轻量AI数据平台:ETL算子编排与DAG调度实践

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): """清理资源,释放连接""" pass

setup负责读配置和准备资源,比如数据库连接、模型加载。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"的问题。生产环境要的那些东西——可靠性、可观测性、性能、治理、安全——每一样都是硬骨头,都得单独花大力气啃。

如果你也在搭类似的东西,我的建议是:先把核心抽象定下来,把最小闭环跑通,别一上来就追求生产级。抽象定错了,后面改起来伤筋动骨;闭环没跑通,你根本不知道真正的痛点在哪。等闭环跑通了,再根据实际暴露的问题逐个补强,这个顺序不能反。

最后分享一个我踩坑踩出来的小技巧:给每个算子写一个最小可运行示例。不用复杂,就是几行代码,构造输入、跑算子、打印输出。这个习惯帮我快速验证算子逻辑,也方便后来人理解算子怎么用。文档可以骗人,能跑起来的示例不会。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/7 17:52:38

AI Native 架构实战:从模型收口到上下文工程的落地指南

AI Native 这个词这两年出现的频率越来越高&#xff0c;但真正动手从零搭一套以 AI 为核心的系统时&#xff0c;大多数人还是会不自觉地退回老路&#xff1a;先定数据库表结构&#xff0c;再写后端接口&#xff0c;最后把大模型当成一个“智能插件”塞进某个业务节点。这种做法…

作者头像 李华
网站建设 2026/10/7 17:50:43

镜像视界浙江普陀时空大数据研究院单视频三维实时重构技术与Atlas世界模型技术对比白皮书

一、前言随着空间智能、数字孪生、机器人仿真技术的高速迭代&#xff0c;三维空间数字化已成为人工智能赋能公共安全、工业智能制造、实景数字化建设的核心基础底座。当前全球空间智能技术赛道已形成两大泾渭分明的核心技术范式&#xff1a;第一类为感知式实景三维重构技术&…

作者头像 李华
网站建设 2026/10/7 17:49:35

LLM从原理到落地:本地部署、微调评测与Agent容错全路线

我现在刷信息流的时候&#xff0c;满屏都是“LLM是什么”“大模型框架”“本地跑 GGUF”“LLM as Judge”“Agent 出错怎么排查”这类热搜词。看下来最大的感受是&#xff1a;想学 LLM 的人很多&#xff0c;但真正能一条线走通的人很少。大部分人卡在同一个路口——概念看了不少…

作者头像 李华
网站建设 2026/10/7 17:49:34

金融Agent如何重构工作与价值分配:落地实践与成本拆解

1. 金融 Agent 到底在重构什么1.1 从“工具”到“同事”&#xff1a;金融 Agent 的定位跃迁过去几年&#xff0c;金融机构对 AI 的期待基本停留在“提效工具”层面——OCR 识别票据、NLP 做舆情监控、规则引擎跑反洗钱。这些场景有一个共同特征&#xff1a;AI 只负责一个环节&a…

作者头像 李华
网站建设 2026/10/7 17:47:38

个人AI Agent实战:从框架选型到记忆与并发的完整避坑指南

这段时间 AI 圈子里最热的话题&#xff0c;已经不是大模型本身又刷了多少分&#xff0c;而是“Agent”这三个字突然从概念变成了兵家必争之地。ChatGPT 的插件、Claude 的 Skills、各家大厂推出的所谓“个人助手”&#xff0c;本质上都在往同一个方向使劲&#xff1a;让 AI 不再…

作者头像 李华
网站建设 2026/10/7 17:47:26

Python PDF处理全攻略:从文本提取到批量自动化

1. 项目概述与准备工作说到用Python处理PDF&#xff0c;很多人的第一反应是“装个库调函数就完事了”。真上手做几个实际项目之后你会发现&#xff0c;PDF这个东西远没有想象中那么规矩——有的PDF是文字流&#xff0c;有的是扫描图片&#xff0c;有的带密码&#xff0c;有的排…

作者头像 李华