应用中的数据流元素较多,并不必然造成高内存或 OOM。实际风险主要取决于同一时段启动的实例数、单实例数据量、流程内部并行度和下游服务承载能力。未运行的数据流通常不会持续占用执行内存,但会增加发布、定时计划、监控和运维管理成本。
数据流的并发可能来自定时重叠、API 集中调用、批量调用其他数据流、失败重试和流程内部并行。并发超过环境的安全容量后,新实例应先排队;如果新增速度长期高于完成速度,队列会越来越长。此时盲目提高并发,往往会把“排队慢”变成“所有任务都变慢或 OOM”。
环境运行容量:代码服务可同时提供的运行工作进程数,主要受 DAGSTER_MAX_WORKER_POOL_SIZE 约束,也受 CPU、内存、数据库和外部接口容量限制。
同一数据流并发数:同一个数据流允许同时处于启动中、进行中或取消中的普通并行实例数,由 DP_FLOW_MAX_CONCURRENCY 控制。它不是整个环境所有数据流的总并发上限,也不应当作所有启动模式的统一限流开关。
单实例内部并发数:一个并行实例内可以同时执行的节点子进程数。它会进一步放大 CPU、内存和数据库连接占用。
预取数量:调度服务提前载入的部分内部调用任务数,由 DP_FLOW_PREFETCH_COUNT 控制。它是缓冲数量,不等于实际运行并发数,也不会增加运行工作进程。
例如,DP_FLOW_MAX_CONCURRENCY=3 表示每个数据流最多有 3 个受该限制的普通并行实例在途。10 条不同的数据流同时触发时,理论需求仍可能达到 30;最终能同时运行多少,还要受工作进程池和环境资源限制。
因此,“超过某个数量后开始排队”的临界点,通常是同一数据流并发上限或工作进程池上限中先达到的一个。如果占用名额的实例长时间不结束、工作进程未正常释放,后续队列就会持续等待;这时应先处理异常实例或服务状态,而不是继续提交任务。
数据流可选择串行或并行启动,具体入口和 InProcess 参数说明参见启动。串行模式在一个执行进程中按依赖顺序执行节点;即使画布上存在多个分支,节点也不会同时运行。并行模式会为已经满足依赖的节点启动子进程,可缩短独立分支的耗时,但资源峰值更高。
画布连线表示数据依赖。下图中,多个前序节点必须完成后,prepared 节点才能开始:

同一节点分出的多条分支如果彼此没有依赖,在并行模式下可能同时运行:

当前版本生成的正式并行任务,会把单实例内部的节点并发上限显式设置为 3。实际同时运行数还会受到可执行节点数量及其他并发限制影响。因此,4 条或更多独立分支不代表它们会全部同时启动;但多个并行实例叠加后,仍可能形成“实例数 × 每实例并行节点数”的资源峰值。
定时计划到达触发时间后会创建一个新实例,并不会等待上一实例完成。如果流程运行时间大于定时间隔,未完成实例就会逐步累积。
预计未完成实例数 ≈ 向上取整(较慢情况下的运行时间 ÷ 定时间隔)
例如,一条流程通常需要 8 分钟,却每分钟触发一次:第 1 个实例尚未完成时,后续 7 个实例已经陆续创建。如果多个实例争用内存、CPU、数据库连接或外部接口,单次耗时还会继续变长,积压速度也会加快。
重叠通常表现为两种情况:达到并发上限的实例停留在【队列中】,实际开始时间越来越晚;或者多个实例同时进入【进行中】,重复执行大查询、连接、排序或 Python 转换,造成内存和数据库压力升高。同一期间被重复处理时,还可能产生重复写入、锁等待或后启动实例覆盖先启动结果。
不要直接按 CPU 核数或任务数量设置并发。建议先用接近生产的数据量运行单个实例,记录以下基线:
P95 运行时长;
相对服务空闲状态增加的峰值内存;
CPU、数据库连接数、临时磁盘和下游接口并发;
是否包含大表连接、排序、去重、透视、窗口计算或 Python DataFrame 转换。
可先按内存估算安全上限:
内存并发上限 = 向下取整(容器内存上限 - 空闲占用 - 20%~30%预留)÷ 单实例 P95 增量峰值内存
再按触发频率估算业务需要的并发:
所需并发 ≈ 高峰期每分钟新增实例数 × P95 运行分钟数
最终并发应取内存、CPU、数据库连接、下游接口和业务需要中的较小值。如果“所需并发”大于“安全上限”,应优先错峰、降低频率、缩短单次耗时、拆分批次或扩容,而不是继续提高并发。
生产前可从 1 个实例开始,按 2、4 等小步增加并发。内存密集型、大表转换或 Python 流程,同一数据流可先从 1~2 开始;普通混合型流程可先从 2~4 开始。以上仅是压测起点,不是固定推荐值。出现内存持续超过 70%~80%、频繁交换内存、数据库连接接近上限、P95 耗时明显增长或下游超时时,应停止提高并发并回退一档。
环境变量属于部署级配置,修改后需要按部署规范重启相关服务。生产环境应由运维人员调整,并保留修改前数值和压测结果。
|
环境变量 |
默认值 |
作用 |
建议 |
|---|---|---|---|
|
|
CPU 核数 |
运行工作进程池的最大容量,通常决定同一代码服务可同时承载的普通运行实例数。池已满时,新实例会等待空闲进程。 |
按前述安全上限设置。CPU 密集型一般不高于可用 CPU 核数;内存密集型应按内存测算值进一步降低。不要仅为清空队列而调大。 |
|
|
2 |
常驻工作进程数量下限,可减少临时创建进程的等待,但会增加基础内存占用。 |
通常保留 1~2;必须小于或等于最大工作进程数。内存紧张时可与最大值一起评估降低。 |
|
|
10 |
按数据流分别限制普通并行实例,启动中、进行中和取消中的实例都会占用名额。串行、同步和调试等入口可能不使用相同的普通运行队列。 |
使用正整数,建议不高于最大工作进程数;设置为 0 会使受该限制的实例无法启动。重型流程宜从 1~2 起测,同时还要限制其他启动入口的提交速率。 |
|
|
64 |
控制同步、批量等部分内部调用任务的预取缓冲数量。 |
一般保持默认值。它不是并发开关;工作进程已满时调大不会加快执行,反而会增加待调度任务的内存占用和管理压力。 |
|
|
30 秒 |
申请空闲工作进程时,单次等待多久后记录超时并重试。 |
通常保持默认值。调大只会延长单次等待,不能增加容量,也不能解决队列不再推进。 |
|
|
未显式配置时为 CPU 核数 |
控制未显式指定上限的多进程执行器中,单实例可同时运行的节点子进程数。 |
当前数据流生成任务已显式配置为 3,因此该变量通常不会改变数据流的单实例内部并发,更不是环境总并发开关。不建议将其作为处理队列积压的手段。 |
DP_WEB_CONCURRENCY、DP_WORKERS_PER_CORE 和 DP_MAX_WORKERS 控制的是 Web/API 服务进程,不是数据流运行实例并发。盲目调大还会增加基础内存占用。DUCKDB_MEMORY_LIMIT 控制单个 DuckDB 连接的内存上限,也不能替代并发控制;如需调整,应同时评估“单连接上限 × 并发连接数”,避免总量超过容器内存。
一组配置必须相互匹配。例如,工作进程池最大值为 4 时,将同一数据流并发设为 10 并不能获得 10 个实际运行进程,只会让更多实例争抢有限资源。反之,如果同一数据流并发设为 2,即使环境仍有空闲工作进程,该数据流的后续普通实例也会继续排队。
当前默认实现中,一个运行工作进程最多复用 100 个实例;当工作进程数量高于最低保有数量时,空闲超过 300 秒的多余进程会被回收。因此,串行实例或节点结束后内存未立刻回落,并不一定代表内存泄漏;应观察多个实例结束后的趋势、空闲回收周期以及容器总内存。若内存跨多个周期持续上升,再结合问题自查与信息采集中的 PID 和 RSS 方法排查。
短时排队是正常的限流结果:运行中的实例完成后,队列中的实例应继续启动。出现以下情况时,不能只按“并发不足”处理:
【进行中】数量始终达到工作进程池上限,且实例仍在正常完成:说明环境已满载。应减少新增任务或按压测结果扩容。
同一数据流已有实例长期处于【启动中】【进行中】或【取消中】:这些实例会持续占用该数据流的并发名额,应先判断是否为长任务、锁等待或异常实例。
【进行中】数量明显低于上限,甚至为 0,但最早排队时间仍不断增加:可能是调度守护进程、代码服务、Redis、运行存储或服务间通信异常,需要运维排查。
日志持续出现申请工作进程超时:说明工作进程池已满或工作进程没有正常释放。调大超时时间只能减少重试频率,不能恢复容量。
实例已经入队,但没有后续的启动记录或运行开始事件:说明问题发生在出队或启动链路,而不是流程内部节点执行缓慢。
排查时建议记录问题时段、最早排队时间、队列数量、启动中/进行中/取消中数量、涉及的数据流、工作进程池配置和服务重启记录。若空闲资源存在但队列超过 5~10 分钟仍完全不推进,应联系运维或技术支持检查队列消费和代码服务状态,不要通过重复提交任务来“唤醒”队列。
定时间隔以接近生产数据量时较慢实例的耗时为依据,并预留数据库波动和资源竞争余量。
多条流程错峰启动,避免集中在整点或同一分钟;批量 API 调用也应限制提交速率。
同一期间、组织或批次不允许重叠时,增加业务唯一键、幂等校验或原子锁。
不要按“每条明细一个被调用实例”展开批量任务,并限制单批参数数量。
先过滤、裁剪字段和聚合,再执行连接、排序或 Python 转换;没有并行要求的分支可增加依赖。
设置超时和重试边界,避免异常实例长期占用工作进程和并发名额。
持续观察“新增速度、完成速度和最早排队时长”。只看当前进行中数量,无法判断队列是否正在失控。
暂停相关定时计划、批量数据流调用和上游 API 调用,先停止新增负载。
按数据流、业务参数和创建时间找出重复及长时间未完成的实例。
评估中断影响后,参考批量终止长时间未完成的实例处理异常实例。
确认调度和代码服务健康、队列能够继续出队后,再以小并发、分批或错峰方式恢复;不要一次性重提全部失败任务。
根据本次峰值重新调整定时间隔、工作进程池、同一数据流并发、批次大小和超时设置。
回到顶部
咨询热线
