
深夜还在处理批量任务的人应该都经历过这种阶段你面对的不是写不出代码而是不敢把代码放出去跑。一个文件压缩任务几十个视频、上千张图片你按顺序一个个处理慢到崩溃用多线程一把梭内存直接被打满拆成几十个子任务丢到队列里跑起来倒是快了可任务到底跑没跑完、哪些失败、卡住的是哪一批你根本说不清楚。你会发现“拆分任务”这件事听起来简单做起来全是不确定性。拆得太粗瓶颈还在单机上拆得太细调度开销超过任务本身拆完之后怎么编排、怎么重试、怎么判断真正完成才是整个流程里最消耗精力的一部分。很多团队一开始只是写个批量处理脚本写到最后硬生生写出了半个分布式调度系统。这篇文章想讲的不是某个新框架而是把这种“拆分 编排”的工程过程理清楚的一种方式。我把它叫做 split dance。理解这个词不用太玄。想象多个执行单元在同一段时间里各自处理自己的碎片既要有分工又不能互相踩踏既要并行推进又要在终点对齐。整个过程有点像一场编排过的舞蹈每一段拆分都有它的节奏。文章主要面向三类读者正在写批量处理、数据同步、定时任务脚本但对并行和任务管理不太放心的开发。想把单机脚本改造成可重试、可观测的小型任务系统的后端工程师。在大模型、视频处理、数据管道场景里处理过大量小任务的 AI 应用开发者。读完之后你可以得到一套容易落地的拆解思路、一个完整的 Python 版最小实现以及生产环境里会遇到的坑和应对方式。1. split dance 到底解决什么问题先想一个日常场景。你有一个日志清洗需求上游给了 200 个压缩文件每个文件可能有几万行日志你需要把它们解压、清洗、过滤、写回数仓。最朴素的写法是 for 循环逐个文件处理。它会慢但不会出事。问题是当文件从 200 变成 20000 时这个循环就成了项目瓶颈。于是你开始改造。你把文件平铺成任务队列开了 16 个线程几万个小任务在队列里飞速流转。这个阶段你会遇到另一类问题任务失败了断点在哪里某几个文件格式不对抛了异常异常把整个 worker 打停了前面的成果要不要重来内存里堆积的中间结果会不会涨到 OOM这时候你会发现真正的复杂度并不是“怎么拆”而是拆完之后拆分出来的单元能不能被独立重试。系统能不能容忍部分成功、部分失败。你能不能在运行过程中看到进度而不是等全部结束之后才面对一堆未知结果。在不丢状态的前提下worker 是否可以随时扩容或退出。split dance 要解决的就是这一整组问题。它是一种工程模式的概括把大任务拆成可独立执行、可追踪、可重试的最小单元再通过调度层把它们组织起来像一支队伍里的舞者一样各自完成动作又在整体上保持同步与一致。这不是某家公司的专利也不是某个开源组件的发明。所有你在生产里见过的任务队列、爬虫调度器、视频转码流水线、多模态数据批处理框架底层都在做这场拆分与编排的 dance。2. split dance 的三个核心角色要理解这种模式先建立一个最小模型。它由三部分构成角色职责对应到工程里的东西Splitter决定把什么拆开、拆得多细文件分片、数据分页、任务分级Worker真正执行一个最小单元线程、进程、容器、远程函数Coordinator管理单元的状态与重试任务表、消息队列、状态机这个模型可以用在任何语言、任何平台上。区别只在于你用线程、进程、Redis 队列还是 K8s Job。2.1 Splitter拆分边界是第一个设计决策拆分不是越细越好。拆分粒度需要考虑两个成本单任务执行成本。调度成本和状态追踪成本。举例来说清洗日志文件时如果把每一行日志都当做一个任务那任务数量会到百万级别调度开销、状态存储开销都不可接受。如果按整个大文件作为一个任务一旦中途出错就可能重跑很久。更合理的拆法可能是每个压缩文件生成一个任务如果单个文件很大再按内部 block 或时间段切片。所以拆分的核心问题是找到一个让“单任务够短、总任务数够可控”的平衡点。经验法则是单个任务执行时间在几秒到几分钟范围通常最好跟踪和重试。2.2 Worker无状态是容器化执行的前提每个最小单元的执行者不应该持有与其他任务相关的状态。如果你需要统计“全局跑了多少条”不要把这个计数放在某个 worker 的内存里而应该由 Coordinator 或外部存储统一维护。这样任何一个 worker 宕掉其他 worker 可以无缝接管它没跑完的任务。这也解释了为什么很多做数据处理的团队最终会走向 Celery、Sidekiq 这类分布式任务框架。它们解决的不是“能不能并发”而是“并发执行时状态到底放在哪”。2.3 Coordinator用一张表看清全局协调者的职责只有一个让每个任务都到达终态。所谓终态就是 success、failed、retry 中的一种。一个高质量的任务系统绝不允许任务停留在“不知道跑没跑”的中间态。工程里最朴素的做法就是一张任务表CREATE TABLE split_task ( id BIGINT AUTO_INCREMENT PRIMARY KEY, task_key VARCHAR(128) NOT NULL UNIQUE, status VARCHAR(16) NOT NULL DEFAULT PENDING, retry_count INT NOT NULL DEFAULT 0, max_retry INT NOT NULL DEFAULT 3, worker_id VARCHAR(64), created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, last_error TEXT ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;task_key 是全流程幂等的关键。它通常由业务主键派生比如文件名加偏移量。有了这张表你随时能回答三个问题当前有多少任务还没跑哪些任务失败次数最多整个批次是不是全部完成3. 环境准备与最小工程骨架现在用一个 Python 示例来跑通整个 split dance 流程。3.1 环境版本本文演示代码不依赖复杂框架建议环境如下Python 3.9 及以上不需要额外安装重依赖标准库即可如果要跑分布式示例需要本机有 Redis并安装 redis-py 库你可以先建一个虚拟环境python3 -m venv splitdance_env source splitdance_env/bin/activate pip install redis3.2 工程目录建议split_dance/ ├── coordinator.py # 调度核心 ├── task_store.py # 任务状态存储抽象 ├── worker.py # 执行器与重试逻辑 ├── splitter.py # 任务拆分逻辑 ├── main.py # 入口示例 └── jobs/ ├── file_job.py # 文件处理业务 └── log_clean_job.py # 日志清洗业务这里先不写太多层。对一个 1000 行以下的内部工具核心代码保持在 3 个文件以内更好维护。4. 核心流程拆解从数据拆分到最终一致把完整流程拆成六个阶段这会比直接贴代码更有用。第一步罗列源头数据。确认你的数据是一组文件、一批数据库记录还是一个需要分段读取的巨型文件。第二步生成 task_key。每个任务单元必须有唯一标识。推荐使用“业务前缀 数据主键”的组合。比如log:20250101:part_3。第三步写入任务表。将 task_key 和任务参数持久化状态置为 PENDING。注意这一步不要和真正执行合并在一起否则一旦执行进程退出任务就丢了。第四步调度拉取。worker 从任务表中捞取 PENDING 状态的记录并通过SELECT ... FOR UPDATE SKIP LOCKED或 Redis 原子操作抢占任务。第五步执行与回调。执行成功后更新状态为 SUCCESS失败时判断 retry_count如果没有超过阈值重置为 PENDING否则置为 FAILED。第六步最终聚合。所有任务都处于 SUCCESS 或 FAILED 终态后由汇总进程读取失败列表决定是否有必要人工介入。很多脚本之所以“一跑就崩、一崩就从头再来”是因为跳过了第一、二、三步直接把业务循环和多线程堆在一起。没有任务清单就没有断点续跑能力。5. split dance 最小实现与代码解读下面进入核心示例代码。5.1 定义任务存储内存版与 Redis 版先把任务存储抽象出来方便从单机内存版切换到分布式。为了篇幅这里同时提供两个版本。# task_store.py import time from dataclasses import dataclass, field from typing import List, Optional dataclass class Task: task_key: str payload: dict status: str PENDING retry_count: int 0 max_retry: int 3 created_at: float field(default_factorytime.time) updated_at: float field(default_factorytime.time) error: Optional[str] None class MemoryTaskStore: 单机演示用进程退出即丢失。 def __init__(self) - None: self._tasks: dict[str, Task] {} def add_task(self, task: Task) - None: if task.task_key in self._tasks: raise ValueError(ftask {task.task_key} 已存在) self._tasks[task.task_key] task def claim_pending(self, worker_id: str, limit: int 1) - List[Task]: claimed: List[Task] [] for task in self._tasks.values(): if task.status PENDING: task.status RUNNING task.updated_at time.time() claimed.append(task) if len(claimed) limit: break return claimed def mark_success(self, task_key: str) - None: task self._tasks[task_key] task.status SUCCESS task.error None task.updated_at time.time() def mark_failure(self, task_key: str, error: str) - None: task self._tasks[task_key] task.error error task.retry_count 1 if task.retry_count task.max_retry: task.status FAILED else: task.status PENDING task.updated_at time.time() def is_all_done(self) - bool: return all(t.status in (SUCCESS, FAILED) for t in self._tasks.values())这段代码用内存字典保存任务。关键在 mark_failure每失败一次就把 retry_count 加一没有超过阈值时重置为 PENDING让 worker 可以重新领取。5.2 定义拆分器与模拟执行器# splitter.py import os from dataclasses import dataclass, field from typing import List dataclass class BatchJob: source_dir: str task_keys: List[str] field(default_factorylist) sleep_seconds: float 0.01 def split(self) - None: 把目录下的所有文件当成最小任务单元。 if not os.path.isdir(self.source_dir): raise FileNotFoundError(f目录不存在: {self.source_dir}) for name in sorted(os.listdir(self.source_dir)): file_path os.path.join(self.source_dir, name) if os.path.isfile(file_path): self.task_keys.append(ffile:{name}) def size(self) - int: return len(self.task_keys)这里展示了最简单的分片逻辑一个文件就是一个任务。更复杂的场景比如需要读取大文件并按字节段切分本质上还是三步计算总长度、按 chunk_size 切点、给每个切点生成 task_key。# worker.py import random from task_store import MemoryTaskStore, Task class DemoWorker: 模拟一个执行器。 def __init__(self, store: MemoryTaskStore, worker_id: str, fail_rate: float 0.2): self.store store self.worker_id worker_id self.fail_rate fail_rate def _process_one(self, task: Task) - None: # 模拟真实处理这里有 20% 概率失败 if random.random() self.fail_rate: raise RuntimeError(f模拟处理失败: {task.task_key}) # 真实项目中在这里调用你的处理函数 # clean(task.payload) print(f[{self.worker_id}] 完成 {task.task_key}) def run(self, batch_size: int 4) - None: while True: tasks self.store.claim_pending(self.worker_id, limitbatch_size) if not tasks: break for task in tasks: try: self._process_one(task) self.store.mark_success(task.task_key) except Exception as exc: # noqa: BLE001 self.store.mark_failure(task.task_key, str(exc))这个 worker 的 while True 写法值得注意它不是一次把所有任务拿走而是每次只领取一批执行完后再拿下一批。这样即使某个批次出现异常代码也不会把整批任务的状态搞乱。5.3 主入口串联全过程# main.py import os import tempfile from splitter import BatchJob from task_store import MemoryTaskStore, Task from worker import DemoWorker def prepare_demo_files(): 创建演示文件。 tmp_dir tempfile.mkdtemp(prefixsplit_dance_demo_) for i in range(20): file_name os.path.join(tmp_dir, flog_{i:03d}.txt) with open(file_name, w, encodingutf-8) as fw: fw.write(fdemo-{i}\n * 10) return tmp_dir def main() - None: store MemoryTaskStore() job BatchJob(source_dirprepare_demo_files()) job.split() print(f拆分出 {job.size()} 个任务) for task_key in job.task_keys: store.add_task( Task( task_keytask_key, payload{path: task_key}, ) ) workers [DemoWorker(store, worker_idfworker-{i}, fail_rate0.1) for i in range(3)] for worker in workers: worker.run() print(f全部终态: {store.is_all_done()}) failed [t.task_key for t in store._tasks.values() if t.status FAILED] print(f失败任务数: {len(failed)}) for task_key in failed[:10]: print(f FAILED - {task_key}) if __name__ __main__: main()运行命令python main.py预期输出结构类似拆分出 20 个任务 [worker-2] 完成 file:log_000.txt [worker-1] 完成 file:log_001.txt ... 全部终态: True 失败任务数: 2 FAILED - file:log_007.txt FAILED - file:log_013.txt因为每个任务最多重试 3 次而 fail_rate 是 0.1所以少量任务在多次重试后仍然失败是合理的。这个输出同样体现了 split dance 的一个重要特征允许部分失败。6. 从单机到分布式Redis 队列版上面内存版能帮你理解核心思路但它有两个硬伤多进程无法共享任务状态进程崩溃后任务丢失。生产环境至少要落到 Redis 或数据库。下面是极简的 Redis 版协调者。# redis_queue_store.py import json import time import uuid from typing import List, Optional import redis class RedisTaskStore: def __init__(self, redis_url: str redis://localhost:6379/0): self.redis redis.Redis.from_url(redis_url, decode_responsesTrue) self.pending_key splitdance:pending self.success_key splitdance:success self.retry_zset splitdance:retry_zset def push_pending(self, task_key: str, payload: dict) - None: self.redis.lpush(self.pending_key, json.dumps({task_key: task_key, payload: payload})) def pop_pending(self, worker_id: str, timeout: int 2) - Optional[str]: raw self.redis.brpoplpush(self.pending_key, fsplitdance:running:{worker_id}, timeouttimeout) return raw def success(self, raw: str) - None: data json.loads(raw) self.redis.sadd(self.success_key, data[task_key]) self.redis.delete(fsplitdance:lock:{data[task_key]}) def failure(self, raw: str, error: str, retry_after_seconds: int 5) - None: data json.loads(raw) score time.time() retry_after_seconds self.redis.zadd(self.retry_zset, {json.dumps(data): score}) def rollback_from_running(self) - int: moved 0 for key in self.redis.scan_iter(matchsplitdance:running:*): while True: raw self.redis.lpop(key) if raw is None: break self.redis.lpush(self.pending_key, raw) moved 1 return moved这个设计并不复杂worker 使用brpoplpush取出任务同时进入自己的 running 队列。成功后写入 success 集合失败后放进延时重试的 zset由另外的恢复逻辑在到点后塞回 pending。宕机后的任务可以通过扫描 running 队列回收到 pending。使用 Redis 后“全部终态”的判断就不适合用所有 worker 完成来判断了因为存在执行中、重试中、运行队列等多种中间状态。正确做法是查询三个指标pending 队列长度。running 队列总长度。retry zset 中 score 小于当前时间的元素个数。三者都归零且 success 数量不为零时才说明一个批次真正结束了。7. 运行结果与效果验证流程跑完验证不能只看“没有任何报错”。在 split dance 模式里正确性等于状态守恒。你可以在单机内存版里加一段校验逻辑def verify(store: MemoryTaskStore, total_expected: int) - None: success_count sum(1 for t in store._tasks.values() if t.status SUCCESS) failed_count sum(1 for t in store._tasks.values() if t.status FAILED) pending_count sum(1 for t in store._tasks.values() if t.status PENDING) running_count sum(1 for t in store._tasks.values() if t.status RUNNING) total_now success_count failed_count pending_count running_count print(f期望任务总数: {total_expected}) print(f当前任务总数: {total_now}, 成功: {success_count}, 失败: {failed_count}, 未终态: {pending_count running_count}) assert total_now total_expected这段断言值得认真对待。分布式执行过程中最常见的 bug 不是某个处理函数报错而是任务在“完成”和“失败”之外又创造了一种隐藏状态既没成功也没失败却永远不会被重新执行。比如 worker 处理任务到一半抛异常但没有捕获比如任务从 Redis 取出后进程被 kill没有回滚比如数据库连接超时导致状态更新失败。这些情况都会造成任务丢失。验证时需要刻意观察一批任务结束后剩余 PENDING / RUNNING 数量是否为零。每个 task_key 是否最多只有一个 worker 在执行。重试的任务是否修改了 payload而不是沿用同一个对象。如果发现日志中有任务被两个 worker 同时处理不要急着去掉重试而要把重点放在业务处理函数是否幂等上。8. 常见问题与排查思路问题现象可能原因排查方式解决方案任务全部卡在 RUNNINGworker异常退出没有回滚运行中任务查看worker退出日志与running队列长度启动恢复任务把超过超时阈值的RUNNING任务重新置为PENDING失败次数很少但任务消失业务代码更新了任务状态异常后无法再触发重试检查每个处理路径是否抛异常统一在Coordinator中维护状态worker只负责业务逻辑重试风暴失败任务过多且重试间隔太短查看重试zset数量和业务日志增加指数退避限制同批任务同时重试数量幂等冲突同一任务被两个worker重复执行写重复数据检查Redis锁或数据库行锁为写操作增加唯一键设置更新任务状态时CAS条件内存增长拆分粒度过细任务对象全量缓存在内存打印任务数改用流式拉取不要一次性把所有任务加载到内存部分成功但整体批次被判定失败汇总逻辑把所有任务当成一个整体事务查看聚合代码明确成功阈值失败任务单独进入补偿流程这里把“幂等冲突”单列出来说因为这是从单机脚本走向分布式任务最常见的坑。单机脚本里同一个任务不会执行两遍分布式系统里网络分区、重试机制、任务恢复都可能让同一逻辑执行两次。只有把业务写成幂等的整个系统才敢放开手脚去自动重试。9. split dance 的工程边界与最佳实践拆分与编排这套模式并不神秘但进入生产环境后有五个原则几乎是无条件的。第一任务清单必须持久化。不要只依赖内存队列把 task_key 和 payload 放进 Redis 或数据库。否则每次发布、宕机、扩容都是一场“从零开始”的噩梦。第二用 task_key 做幂等。每个任务在执行前先记录自己的 task_key。无论 Redis 锁、数据库唯一约束还是消息队列的消息去重task_key 都是保障一致性的锚点。第三重试一定要有退避。失败后立刻重试只会让问题雪上加霜。推荐指数退避第一次失败 5 秒后重试第二次 25 秒第三次 125 秒。超过最大次数后进入 FAILED 状态等待人工或告警处理。第四日志要按 task_key 贯穿。业务执行和任务状态分开打日志每条执行日志都带上 task_key。这样排查问题时你从任务表里看到失败的 key再去日志系统里 grep 同一个 key就能还原整个执行链路。第五灰度与预览先行。新加一个文件处理类型时先拆分出 1% 的任务试运行观察成功率、耗时、输出格式再扩大到全量。不要在一次发布里把 20 万任务直接丢进新代码。团队协作上如果多个人同时开发同一套任务系统必须约定任务参数 schema。建议使用 JSON 并做版本号管理{ version: 1, task_type: log_clean, source_path: /data/input/2025/01/01/log_001.txt, output_table: clean_logs }不要让 worker 直接依赖某个字段名而是通过 task_type 路由到不同处理函数。这样后续新增任务类型时只需要注册新的 handler不需要改动调度核心。如果你在做 AI 应用这种按任务类型路由的方式也适合接模型调用把图片、文本、音视频各自的处理封装成不同 handler再由同一套调度系统统一排队重试能帮你在模型限流、任务积压时保住数据完整性。10. 实际项目中的演进路径最后给出一条建议的落地路径避免一上来就上重框架。第一阶段单机脚本加断点续跑。把任务改成“先写任务表再执行执行完更新状态”。哪怕只是单进程顺序执行也已经建立起任务完整性意识。第二阶段线程池/进程池执行。设计成一个消费任务列表的 worker通过concurrent.futures控制并发数。引入失败重试和终态判断这基本就是单机版 split dance。第三阶段进程级隔离。多个 worker 进程同时跑任务表和状态落到 Redis 或 MySQL。验证并发量上去后任务不丢、不重复。第四阶段再判断是否需要引入 Celery、Temporal、Argo Workflows 这类成熟框架。不是所有场景都需要框架如果只是每天跑几百个文件前面三个阶段已经足够。如果任务量达到数万级别且有复杂的超时、重试、定时、人工审批链路成熟框架的收益会大很多。这里想强调一个容易被忽略的点框架不是银弹。直接引入 Celery 并不能让你自动获得“合理的任务粒度”和“幂等的业务逻辑”。相反框架提供了更多旋钮配置不当会让系统更难排查。先理解这场拆分之舞的每个节拍再去用框架表达比反过来学要快得多。如果你正在推进类似的数据批处理、AI 推理批处理、文件转码、爬虫抓取任务可以先拿本文的单机内存版代码跑通全流程再逐步替换成 Redis 或数据库版本。把示例中的BatchJob.split()换成你的真实数据源把DemoWorker._process_one()换成你的真实业务函数很快就能感受到这套思路带来的确定性。建议收藏这份代码骨架后续新增批量任务时直接作为模板使用。