十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

任务依赖实战:从“双数组先做”到自动化状态检查

任务依赖实战:从“双数组先做”到自动化状态检查 如果团队里出现一句“双数组任务先做34 号完成了我们再动手”这表面上是排期沟通实际上已经点出了一个多任务并行开发中非常核心的问题任务之间存在前置依赖而依赖关系如果不拆清楚后面要么有人空等要么两个任务各改各的合并时才发现冲突。这篇文章不打算讲某个具体开源项目而是把这句话里涉及的两类任务当作一个典型场景来分析。我们来看“先做双数组”和“等 34 号完成后再启动”这两类任务在真实开发流程中分别意味着什么、要用什么方式管理依赖、怎么做批量任务验证、怎么避免上游延期把下游一起拖死。文章会给出可以直接套用的任务拆分方案、状态检查脚本、测试流程和问题排查清单。1. 任务本质拆解这句话到底在说什么“双数组先做”和“34 做了我们就做”两句话信息量很大拆开看是三层意思。第一层双数组任务是第一个要启动的工作项。“双数组”本身在不同语境下差别很大。可能是需要对两个无序数组合并去重可能是双数组 Trie 的构建也可能是二分查找、双指针排序这一类需要同时操作两个数组的算法任务。无论具体是哪种它在技术实现上都有共同点要定义清楚两个数组的输入格式、输出格式、数据规模、边界条件。第二层34 号任务是一个前置阻塞项。“34 做了我们就做”说明 34 号任务是当前任务的上游。上游没完成下游能做的事情是有限的但有限不等于完全不能做。下游可以先把接口定义好、把测试数据准备好、把处理流程的骨架搭好只等上游输出一个确定的结果就能立刻接入联调。第三层团队里存在并行协作和等待关系。如果团队里同时有“先做 A”和“等 B 做完再执行 C”两条指令那么这些任务之间构成了一个典型的有向无环图DAG。A 任务可以尽快启动C 任务依赖 B 完成。如果没有把这种依赖关系显式管理起来完全靠口头同步最容易出现的情况是下游团队成员反复问“34 好了吗”或者上游以为自己交付了但下游拿到的格式不是预期格式又得返工。一句话总结这句话的本质是任务编排问题不是单纯的技术问题。要想两边都不卡住需要把任务拆解细化把依赖从口头约定变成自动化状态检查。2. 双数组任务的启动方式先定契约再写实现既然“双数组先做”那么第一个动作不是打开编辑器写代码而是先确定输入输出契约。以一个常见的双数组合并去重任务为例。上游给两个数组下游需要把两个数组合并后按升序输出并且要去掉重复值。这类任务听起来简单但一旦真正落地第一步就要讨论清楚几个细节。数组是内存中的 List还是来自文件、数据库、消息队列数组长度量级是百、万、百万还是亿元素类型是整数、字符串还是对象输出顺序是否敏感去重标准是什么相同对象的比较字段是哪个这些问题没定清楚就动手写出来的代码很可能要推翻。经验做法是在项目目录里先维护一份contract.md或者接口定义文件把输入输出字段、类型、示例全部列出来。还没有真实上游数据时就先造一份 mock 数据用同样的格式来驱动开发。一份 mock 输入数据的结构可以是这样的{ task_id: task_double_array_001, input: { array_a: [3, 1, 4, 1, 5], array_b: [9, 2, 6, 5, 3] }, config: { sort: true, deduplicate: true } }对应的处理函数只需要保证接收类似结构的输入返回统一结构的输出。只要这个结构定了后续 34 号任务交付什么格式都不影响这边的主体代码。下面是双数组任务的核心处理部分代码本身不复杂重点在于边界处理from typing import List def merge_two_arrays( array_a: List[int], array_b: List[int], sort: bool True, deduplicate: bool True ) - List[int]: if not array_a and not array_b: return [] if deduplicate: merged list(set(array_a) | set(array_b)) else: merged array_a array_b if sort: merged.sort() return merged if __name__ __main__: sample { array_a: [3, 1, 4, 1, 5], array_b: [9, 2, 6, 5, 3] } result merge_two_arrays(**sample) print(result)这段代码不是银弹它想说明的是双数组任务在正式进入“批量处理”之前需要先用一个最小样例把数据链路跑通。跑通之后后面无论接 34 号任务的数据还是接其他上游数据都只是格式适配的问题。如果“双数组”指的是双数组 Trie 这类数据结构思路也完全一致。开工前先明确构建原料、查询模式、内存预算、是否支持动态插入然后先写最小可运行版本再压测。数据结构类任务最容易踩的坑不是不会写而是没有先确认数据规模就盲目追求高级实现。3. 34 号前置任务的状态管理从口头询问到可查询接口“34 做了我们就做”这句话最大的风险在于依赖关系靠人肉记忆。如果 34 号任务是某一次代码评审的编号、某个 Bug 的修复单号、某一次数据迁移的批次号或者某个接口的上线编号那么下游团队需要随时知道它当前处于什么状态。理想情况是34 号任务完成后会产生一个可被程序感知的状态变更。这种变更通常有四种载体。代码仓库主分支上出现了某个特定提交。CI 流水线执行成功并产出了可下载的构建物。某个数据库表或状态文件被更新。某个接口返回了 ready 状态。具体用哪一种取决于公司的技术栈。但不管哪种下游团队都不应该靠“问一嘴”来获取状态而应该用一个自动化脚本去检查。假设 34 号任务的完成标志是远端 Git 仓库出现了一个 tagrelease-task-34。那么下游任务启动前可以先执行状态检查#!/usr/bin/env bash TAG_NAMErelease-task-34 REMOTEorigin if git ls-remote --tags $REMOTE $TAG_NAME | grep -q $TAG_NAME; then echo upstream task 34 is done, ready to start downstream task exit 0 else echo upstream task 34 is not ready, waiting... exit 1 fi如果检查不通过脚本返回非 0 退出码。这个退出码可以直接被 CI 系统识别让下游流水线处于阻塞状态而不是直接失败。上游 tag 一打出来下一次轮询或者 webhook 触发之后下游流水线就会自动继续执行。这样一来“34 做了我们就做”就从一句口头承诺变成了可自动感知的流水线门禁。如果 34 号任务的完成标志是接口状态则检查逻辑差不多import requests import time UPSTREAM_STATUS_URL http://your-ci-server/api/tasks/34/status CHECK_INTERVAL_SECONDS 60 TIMEOUT_SECONDS 3600 start_time time.time() while time.time() - start_time TIMEOUT_SECONDS: try: response requests.get(UPSTREAM_STATUS_URL, timeout10) data response.json() status data.get(status) if status success: print(upstream task 34 finished, start downstream task) break else: print(fupstream status is {status}, check again after {CHECK_INTERVAL_SECONDS}s) except Exception as exc: print(fcheck failed: {exc}) time.sleep(CHECK_INTERVAL_SECONDS) else: raise RuntimeError(timeout waiting for upstream task 34)需要注意这里的状态轮询必须设置超时时间和失败重试逻辑。否则 CI 任务会因为网络抖动或上游任务临时挂起而无限等待挤占流水线资源。比较稳妥的实践是轮询间隔设 30 到 60 秒超时时间设 1 到 4 小时超过时间直接给出告警由负责人确认上游是否出现阻塞。4. 批量任务处理双数组与 34 号产出对接后的自动化下游任务在拿到 34 号产出的数据之后面对的很可能不是单条输入而是一批数据。比如 34 号任务产出了一份包含 1000 组数组对的文件下游需要用双数组处理逻辑逐一处理并汇总输出。这就是批量任务的典型场景。批量任务处理最忌讳的是在循环里直接打印日志、不做失败隔离、不做中间结果持久化。一个输入出错整个批量任务中断重新启动以后又要从头跑。这种情况在数据量大时非常浪费时间。更合理的做法是设计一个具备三个能力的批量脚本单条失败不影响整体、处理进度可恢复、结构化的输入和输出文件。import json import os import logging from pathlib import Path logging.basicConfig(levellogging.INFO, format%(asctime)s - %(levelname)s - %(message)s) logger logging.getLogger(__name__) INPUT_JSONL ./inputs/task_34_output.jsonl OUTPUT_JSONL ./outputs/double_array_results.jsonl CHECKPOINT_FILE ./outputs/checkpoint.txt def process_one_line(line: str) - dict: record json.loads(line) array_a record[array_a] array_b record[array_b] result sorted(set(array_a) | set(array_b)) return {task_id: record.get(task_id), result: result} def load_checkpoint() - int: if os.path.exists(CHECKPOINT_FILE): with open(CHECKPOINT_FILE, r, encodingutf-8) as f: return int(f.read().strip()) return 0 def save_checkpoint(line_number: int) - None: with open(CHECKPOINT_FILE, w, encodingutf-8) as f: f.write(str(line_number)) def main() - None: Path(OUTPUT_JSONL).parent.mkdir(parentsTrue, exist_okTrue) start_line load_checkpoint() logger.info(batch task start from line %s, start_line) with open(INPUT_JSONL, r, encodingutf-8) as fin, \ open(OUTPUT_JSONL, a, encodingutf-8) as fout: for line_number, line in enumerate(fin, start1): if line_number start_line: continue try: result process_one_line(line) fout.write(json.dumps(result, ensure_asciiFalse) \n) fout.flush() save_checkpoint(line_number) except Exception as exc: logger.warning(line %s process failed: %s, line_number, exc) continue logger.info(batch task finished) if __name__ __main__: main()这个脚本的核心思路是断点续跑。处理成功一行就把当前行号落到 checkpoint 文件里。即使中途因为内存问题、机器重启、数据异常中断再次启动时也能直接从上次位置继续。使用 JSONL 而不是一个大 JSON 数组也是为了让每一行都能被独立解析不会因为一条脏数据导致整个文件读不出来。如果输入不是 JSONL 而是普通文本行或者上游产出的是一整个目录的多份文件只需要把process_one_line换成process_one_file整体框架不变。批量任务的关键从来不是单条逻辑写得多精巧而是失败恢复和数据持久化是否可靠。5. 功能测试与效果验证判断两个任务是否真的完成在投入联调和批量任务之前需要先定义清楚“完成”的校验标准。测试不只是为了证明代码能跑通更是为了给前面那两句话一个明确答复双数组任务到底做到什么程度可以算完成34 号任务到底做到什么程度下游才能启动。双数组任务建议按下面几个维度来验证。基础正确性。输入两个有序数组合并后的结果是否仍然有序。去重正确性。两个数组内部有重复值时输出是否去掉了重复元素。空数组和单元素数组。边界输入情况下程序是否崩溃。大数组性能。数组长度达到万、十万、百万量级时执行耗时是否在可接受范围内。输出可重复性。同一份输入执行多次结果是否稳定一致。内存占用。处理超大数组时是否因为频繁复制导致内存飙升。测试代码可以很朴素import time import random test_cases [ ([], [], []), ([1], [], [1]), ([1, 2, 3], [1, 2, 3], [1, 2, 3]), ([3, 1, 4], [9, 2, 6], [1, 2, 3, 4, 6, 9]) ] def run_basic_tests(): for a, b, expected in test_cases: result merge_two_arrays(a, b) assert result expected, fcase failed: {a}, {b}, got {result}, expected {expected} print(basic tests passed) def run_perf_test(array_len: int 100000): a [random.randint(0, 100000) for _ in range(array_len)] b [random.randint(0, 100000) for _ in range(array_len)] start time.time() merge_two_arrays(a, b) cost time.time() - start print(farray length {array_len}, cost {cost:.2f}s) if __name__ __main__: run_basic_tests() run_perf_test()34 号任务作为上游同样需要一份“可启动下游”的验收清单。至少应该满足34 号任务完成后产出的数据文件是否存在、文件内部格式是否与约定一致、抽样数据是否通过校验等。如果上游产出的是接口服务则需要检查接口是否能稳定响应、响应时间是否达标、鉴权是否可配置。这里给出一个上游产出物的校验脚本思路它不针对特定项目但可以作为通用模板import json from pathlib import Path def validate_upstream_file(file_path: str) - bool: data_file Path(file_path) if not data_file.exists(): print(upstream output file not exists) return False with open(data_file, r, encodingutf-8) as f: for line_number, line in enumerate(f, start1): try: record json.loads(line) if array_a not in record or array_b not in record: print(fline {line_number} missing required field) return False except json.JSONDecodeError: print(fline {line_number} is not valid json) return False print(upstream output file validation passed) return True if __name__ __main__: ok validate_upstream_file(./outputs/upstream_task_34.jsonl) exit(0 if ok else 1)实际项目里这个校验脚本应该作为下游流水线的第一个阶段。上游没有产出或产出物不合法时下游直接终止并返回一个明确的原因而不是等到跑批跑到一半才发现数据有问题。6. 接口 API 与任务状态对接如何把依赖自动化如果你的团队开发环境里34 号任务是某个服务集群上的异步任务那么状态对接就需要通过 API 来做。这里需要强调一点不要在业务代码里到处写死“等待 34 号任务”的逻辑。更好的方式是把状态检查抽成一个独立的依赖服务或者一个独立函数同时预留重试和超时。这样即使任务编号变化比如以后出现“45 做了再做”只需要修改配置不需要重写逻辑。一个可以放到配置中心的示例{ dependency_task_id: 34, upstream_status_api: http://service.internal/api/v1/tasks/{task_id}/status, check_interval_seconds: 30, timeout_seconds: 7200, on_success_hook: http://pipeline.internal/api/v1/trigger/double_array_batch }下游脚本读取这个配置轮询上游任务状态当状态为 success 时调用 on_success_hook 触发自己的批量任务。如果自研成本高也可以直接使用现成的 CI 或工作流引擎来配置依赖关系。GitLab CI 的needs关键字、Jenkins 的build触发条件、阿里的流水线编排、GitHub Actions 的workflow_run事件都属于“上游完成后再跑下游”的成熟方案。比较推荐的做法是如果团队已经使用了某种 CI 平台优先用平台自带的依赖编排能力而不是自己写轮询脚本。只有在平台能力覆盖不到或者需要跨系统协调时才自己维护状态查询脚本。下面给一个 GitHub Actions 风格的 YAML 示例用来表达“下游任务等上游成功后再触发”的编排思路。不同平台的字段差异较大使用时需要按团队实际平台调整name: downstream-double-array on: workflow_run: workflows: [upstream-34-task] types: - completed jobs: check-upstream-status: runs-on: ubuntu-latest steps: - name: check upstream result run: | echo upstream workflow completed echo if you need to check its conclusion, use GitHub Actions API如果上游 workflow 实际上是失败的下游也应该根据结论自动跳过。这个在不同平台实现方式不一样但核心判断逻辑是状态 success 才继续状态 failure 则发告警状态 pending 或 running 则继续等待。7. 资源占用与性能观察别等任务跑完了才发现内存不够“双数组先做”听着简单但如果数组规模大资源占用也是实打实的问题。尤其是两个数组都在内存里合并时又产生一个新的数组内存消耗很容易翻倍。需要掌握几个观察手段。在 Python 里可以用tracemalloc来统计内存占用import tracemalloc import random array_len 500000 a [random.randint(0, 100000) for _ in range(array_len)] b [random.randint(0, 100000) for _ in range(array_len)] tracemalloc.start() result merge_two_arrays(a, b) current, peak tracemalloc.get_traced_memory() tracemalloc.stop() print(fcurrent memory: {current / 1024 / 1024:.2f} MB) print(fpeak memory: {peak / 1024 / 1024:.2f} MB)提高性能的方向取决于双数组任务的具体语义。如果是两个有序数组合并完全不需要set后再sort用双指针归并就是 O(n) 的时间复杂度。如果数组量级很大可以考虑用numpy向量化操作。但引入 numpy 之前要先评估运行环境是否具备安装条件。from typing import List def merge_sorted_arrays(array_a: List[int], array_b: List[int]) - List[int]: i 0 j 0 result [] while i len(array_a) and j len(array_b): if array_a[i] array_b[j]: result.append(array_a[i]) i 1 else: result.append(array_b[j]) j 1 if i len(array_a): result.extend(array_a[i:]) if j len(array_b): result.extend(array_b[j:]) return result这个版本只适用于输入有序的情况。不要盲目套用先确认上游数据是否有序再决定算法。数据规模只有几千的时候O(n log n) 和 O(n) 差别不大直接写简单方案更不容易出错。到了百万级才需要认真考虑算法复杂度和内存布局。34 号上游任务如果本身是长时间运行的批处理或模型推理任务还需要关注它对 CPU、GPU、内存和磁盘的占用情况。比如在 Linux 服务端执行长时间任务时可以用以下命令观察前后对比top -b -n 1 | head -20free -hdf -h如果 34 号任务跑在一台 8 G 内存的机器上上游任务执行过程占用接近 7 G此时又强行启动双数组批处理进程两个任务同时挤在同一台机器上大概率会出现 OOM。观察资源占用不只是开发阶段的事生产调度阶段更重要。很多下游任务失败不是代码逻辑错而是机器资源不够或者两个任务在错误的时间重叠了。8. 常见问题与排查方法从空等到数据错乱两个任务相互依赖时最常遇到的问题就集中在几个点上上游任务状态更新缺失、批处理脚本中断无法恢复、输出数据格式不兼容、任务执行过程中资源不够。问题现象可能原因排查方式解决方案下游一直空等无法启动上游任务没有更新状态状态检查脚本没有触发查看上游任务日志和 CI 状态给上游添加完成回调缩短轮询间隔增加手动重试入口下游启动后立刻报错34 号任务产出物还没生成或路径不对检查路径下文件是否存在文件大小是否为 0校验产出物文件后再触发下游批量任务跑一半中断进程 OOM、机器重启、脚本异常退出查看进程日志和系统 dmesg引入 checkpoint 断点续跑增加失败重试双数组结果与预期不一致输入数组顺序未确认、去重字段判断标准不一致对比输入样例和输出样例先固化 mock 数据与预期结果再与上游对齐格式接口轮询导致服务压力大轮询间隔太短请求并发太高查看上游服务 QPS 和日志改为更长的轮询间隔或者用 webhook 推送替代轮询上游与下游并行执行修改同一批文件任务依赖只做了部分编排没有加锁检查两个任务的工作目录是否重叠分目录管理输入输出设定文件锁或使用独立工作区实际排查这类问题一般先看两个时间点。第一个是上游完结时间点确认上游是否真的已经产出完整文件第二个是下游启动时间点确认下游读取到的文件是完整版本。中间任何一个环节出现时间差都可能读到半截文件直接导致解析失败。批处理脚本如果卡住不要只靠肉眼看终端。要检查输出文件行数是否还在增长wc -l ./outputs/double_array_results.jsonl多次执行该命令如果行数一直在增加说明任务没死只是慢。如果行数长时间不变说明任务可能已经阻塞。此时需要查看进程状态和日志不能盲目重启进程。强行重启可能丢失一部分已处理结果虽然 checkpoint 机制能够恢复但最好的做法是先确认进程是否还活着。9. 最佳实践与使用建议把依赖关系变成工程能力这类“上游做完我再做”的流程在团队协作中会反复出现。与其每次口头对齐不如把下面几件事固化到日常研发流程里。首先任务启动前先在文档里画清楚依赖关系。不需要多复杂的图一张表格就够了任务前置任务启动条件产出物验证方式双数组合并任务无开一个任务分支即可启动合并后的数组集合文件单测、mock 数据双数组批量处理任务34 号任务34 号任务状态为 success批量结果 JSONL抽样校验、行数校验这张表的作用是让每个人都知道自己什么时候能动、什么时候要等、等的东西是什么。任务挂在看板上以后状态就变成了一种可查询的资源。其次主流程代码与适配代码要分离。双数组处理逻辑只关心“输入数组 A 和数组 B输出结果数组”至于数组来自 34 号任务还是来自手工上传应该在入口层做适配。很多人写代码容易把上游字段直接透传到核心逻辑里比如上游字段叫listA核心逻辑里也写listA一旦上游改字段名核心逻辑也跟着改。这类代码耦合会让协作成本越来越高。更合理的目录规划大致如下project/ ├── adapters/ # 针对 34 号任务的特殊字段适配层 ├── core/ # 双数组处理核心逻辑不依赖任何外部字段 ├── inputs/ # 上游输入数据 ├── outputs/ # 输出结果与日志 ├── tests/ # 单测与集成测试 └── scripts/ # 状态检查、批量启动、checkpoint 工具再次要重视 mock 先行。上游还没完成时下游最应该做的是构造一份和约定格式完全一致的 mock 数据提前把下游流程全部跑通。等上游真实数据接入时只需要把 mock 数据源切换成真实数据源一般最多只需要处理几个字段名不完全一致的小问题不会出现流程性的阻塞。另外批处理任务的日志规范一点。每个文件或每一行处理完成后输出带行号或任务 ID 的日志方便出问题时快速定位。不要在整个批处理结束以后才输出一条总日志这样中间出错连位置都找不到。最后也是非常重要的一点如果项目中涉及的是带版权、肖像权或隐私的数据比如名字叫“双数组”但实际里面存放的是敏感信息或者 34 号任务是人脸、声音、文档等敏感数据的处理任务那么分布式协作、数据导出、共享链路的每一个环节都要先确认授权与合法性。技术排期只是一部分数据合规边界必须前置。不要为了让任务快速流转就把未经处理的用户数据直接放在共享目录里或用明文接口传递。这部分一定要在任务启动之前单独确认。10. 总结与下一步这一次先动手做哪件事回到开头那一句“双数组先做34 做了我们就做”最值得尝试的点是不要把它当成一次口头排期而是把它当作一次小规模的依赖编排演练。要在项目里先落地的内容是把任务拆到可以验证的粒度。上游任务 34 需要一个可查询的完成状态双数组任务需要一份 mock 输入和一份明确的输出样例。然后写一个简单的状态检查脚本把等待过程自动化。脚本不复杂几十行代码就够用。最后跑一个包含 10 条左右输入的批量任务验证中断恢复和输出格式是否稳定。最容易踩的坑有三个。第一个是上下游对“完成”的定义不一致上游认为代码合入就算完成下游需要的是某个数据文件生成。第二个是批量处理没有断点续跑机制跑一半挂了以后只能从头再来。第三个是输入输出格式没有提前固化导致 34 号任务真的做完时下游还在适配字段名称白白把并行开发的红利消耗掉。一旦把双数组任务跑通把 34 号任务的状态检查接好这套依赖管理方式可以继续复用到后面的 “45 号任务做完再做”“50 号接口发布后再批量跑”等更多场景。本质不变都是把消息传递变成可验证的接口和文件把人的记忆变成自动化的状态判断。如果所在团队还没有统一的任务看板和 CI 依赖编排可以先从一份任务依赖表和一个状态查询脚本起步。这两样东西不需要引入复杂的平台却能在最短时间内消除“他到底做完没有”这个最常见的协作黑洞。
返回列表