You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何确保批量插入最后一批数据时仅生成一次计算进程?

解决并发场景下仅启动一次计算进程的方案

针对你遇到的并发请求重复启动计算进程的问题,提供以下几种可行的解决方案:

方案一:利用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对应的计算进程只能被启动一次,后续请求会直接复用已存在的进程。

方案二:数据库层面加排他锁(适合强一致性场景)

通过数据库事务加排他锁,将「判断记录是否完成」和「启动进程的标记」变成原子操作,避免并发冲突:

  1. 先给groups表新增一个processing_started字段(布尔类型,默认false),用于标记计算进程是否已启动。
  2. 修改代码,在事务中执行判断和标记:
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表来做全局锁,保证跨节点的唯一性:

  1. 在应用启动时创建ETS表(例如在application.ex的start/2方法中):
def start(_type, _args) do
  # 创建全局可见的ETS表,用于记录已启动计算进程的group_id
  :ets.new(:processing_groups, [:named_table, :public, :set])

  # 其他启动逻辑...
end
  1. 修改业务代码:
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 14:20:50