如何确保批量插入最后一批数据时仅生成一次计算进程?
解决并发场景下仅启动一次计算进程的方案
针对你遇到的并发请求重复启动计算进程的问题,提供以下几种可行的解决方案:
方案一:利用Elixir进程唯一注册特性
Elixir中,当给进程指定name参数启动时,如果同名进程已经存在,start_link会返回{:error, {:already_started, pid}}。我们可以捕获这个错误,避免重复启动:
def bulk_create(conn, %{"records" => records_params}) do with {rows_affected, nil} <- Api.create_records(records_params) do group_id = records_params |> first() |> Map.get("group_id") group = Api.get_group!(group_id) total_records_saved = Api.count_records_by_group_id(group.id) finished = group.total_records == total_records_saved pid = if finished do process_name = String.to_atom("PID.#{group.id}") case Processor.start_link(%{group_id: group.id, name: process_name}) do {:ok, pid} -> pid |> :erlang.pid_to_list() |> Enum.slice(1..-2) |> List.to_string() {:error, {:already_started, existing_pid}} -> existing_pid |> :erlang.pid_to_list() |> Enum.slice(1..-2) |> List.to_string() end else nil end conn |> put_status(:created) |> render("bulk_create.json", rows_affected: rows_affected, finished: finished, pid: pid ) end end
这个方案的核心是依赖Elixir的进程注册机制,确保同一个group_id对应的计算进程只能被启动一次,后续请求会直接复用已存在的进程。
方案二:数据库层面加排他锁(适合强一致性场景)
通过数据库事务加排他锁,将「判断记录是否完成」和「启动进程的标记」变成原子操作,避免并发冲突:
- 先给
groups表新增一个processing_started字段(布尔类型,默认false),用于标记计算进程是否已启动。 - 修改代码,在事务中执行判断和标记:
def bulk_create(conn, %{"records" => records_params}) do with {rows_affected, nil} <- Api.create_records(records_params) do group_id = records_params |> first() |> Map.get("group_id") {pid, finished} = Repo.transaction(fn -> # 加FOR UPDATE锁,确保同一时间只有一个请求能操作该group group = Repo.one!(from g in Group, where: g.id == ^group_id, lock: "FOR UPDATE") total_records_saved = Api.count_records_by_group_id(group.id) finished = group.total_records == total_records_saved pid = if finished and not group.processing_started do # 标记为已启动 Repo.update!(Group.changeset(group, %{processing_started: true})) # 启动进程 {:ok, pid} = Processor.start_link(%{group_id: group.id, name: String.to_atom("PID.#{group.id}")}) pid |> :erlang.pid_to_list() |> Enum.slice(1..-2) |> List.to_string() else nil end {pid, finished} end) conn |> put_status(:created) |> render("bulk_create.json", rows_affected: rows_affected, finished: finished, pid: pid ) end end
这里通过PostgreSQL的FOR UPDATE锁,确保同一时间只有一个请求能进入事务逻辑,要么启动进程并标记状态,要么看到已标记就跳过启动,从根源避免重复。
方案三:用ETS表实现分布式锁(适合分布式部署场景)
如果你的服务是分布式部署的,可以用ETS表来做全局锁,保证跨节点的唯一性:
- 在应用启动时创建ETS表(例如在
application.ex的start/2方法中):
def start(_type, _args) do # 创建全局可见的ETS表,用于记录已启动计算进程的group_id :ets.new(:processing_groups, [:named_table, :public, :set]) # 其他启动逻辑... end
- 修改业务代码:
def bulk_create(conn, %{"records" => records_params}) do with {rows_affected, nil} <- Api.create_records(records_params) do group_id = records_params |> first() |> Map.get("group_id") group = Api.get_group!(group_id) total_records_saved = Api.count_records_by_group_id(group.id) finished = group.total_records == total_records_saved pid = if finished do process_name = String.to_atom("PID.#{group.id}") # 原子操作:只有当group_id不在表中时才插入成功 case :ets.insert_new(:processing_groups, {group_id, true}) do true -> {:ok, pid} = Processor.start_link(%{group_id: group.id, name: process_name}) pid |> :erlang.pid_to_list() |> Enum.slice(1..-2) |> List.to_string() false -> # 已存在进程,返回已有PID existing_pid = Process.whereis(process_name) existing_pid |> :erlang.pid_to_list() |> Enum.slice(1..-2) |> List.to_string() end else nil end conn |> put_status(:created) |> render("bulk_create.json", rows_affected: rows_affected, finished: finished, pid: pid ) end end
ets.insert_new是原子操作,能保证跨节点只有第一个请求能插入成功并启动进程,其他请求会直接复用已有的进程。
内容的提问来源于stack exchange,提问作者Pablo Moltedo
相关产品推荐
相关产品推荐

