如何基于OTP实现500个一次性任务的分批/串行延迟执行?
基于OTP的异步任务调度实现
首先明确::timer.sleep/1确实会阻塞当前进程,如果在主线程调用,整个应用都会卡住,绝对不能这么用。下面是两种方式的实现,完全依赖OTP组件,无需额外依赖:
通用工具函数
先实现随机延迟生成,满足可选需求:
# 生成1-4秒之间的随机毫秒数 defp random_delay(), do: Enum.random(1000..4000)
方式b:逐个执行任务(任务间延迟)
核心思路是异步启动任务,然后通过Process.send_after/3触发下一个任务的执行,全程不阻塞主线程:
脚本/独立进程版本
def run_tasks_one_by_one(tasks) do # 启动调度逻辑 send(self(), {:run_next, tasks}) # 消息循环处理调度 loop() end defp loop() do receive do # 任务列表为空时结束 {:run_next, []} -> :ok {:run_next, [task | remaining_tasks]} -> # 异步启动当前任务,无需等待结果 Task.start(task) # 延迟后触发下一个任务 delay = random_delay() Process.send_after(self(), {:run_next, remaining_tasks}, delay) # 继续等待消息 loop() end end
GenServer集成版本(如果在应用进程中使用)
如果是在GenServer里实现,只需添加对应的消息处理:
def init(_args) do # 假设任务列表存在state里,或者外部传入后触发 {:ok, %{tasks: all_tasks}} end def handle_info({:run_next, []}, state) do {:noreply, state} end def handle_info({:run_next, [task | rest]}, state) do Task.start(task) delay = random_delay() Process.send_after(self(), {:run_next, rest}, delay) {:noreply, %{state | tasks: rest}} end
方式a:分批次执行(批次间延迟)
先把任务分成8-10个一组的批次,再按批次异步启动,批次间用延迟分隔:
分批工具函数
# 将任务列表拆分为8-10个一组的批次 defp split_into_batches(tasks) do # 随机每组8-10个任务 batch_size = Enum.random(8..10) Enum.chunk_every(tasks, batch_size) end
批次调度实现
def run_in_batches(tasks) do batches = split_into_batches(tasks) send(self(), {:run_batch, batches}) loop() end defp loop() do receive do {:run_batch, []} -> :ok {:run_batch, [batch | remaining_batches]} -> # 异步启动当前批次的所有任务 Enum.each(batch, &Task.start/1) # 延迟后触发下一批次 delay = random_delay() Process.send_after(self(), {:run_batch, remaining_batches}, delay) loop() end end
关键注意点
- 用
Task.start/1启动任务:任务会在独立进程中运行,不会阻塞调度逻辑 - 用
Process.send_after/3实现延迟:这是异步操作,不会占用当前进程的执行时间 - 随机延迟:
Enum.random(1000..4000)可以直接生成指定范围的随机毫秒数,简单直观 - 如果需要监控任务状态,可以改用
Task.Supervisor.start_child/2,配合监督树管理任务生命周期
内容的提问来源于stack exchange,提问作者Marco C. Stewart
相关产品推荐
相关产品推荐

