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

Elixir Umbrella App:Oban通知无法跨Repo更新产品缓存

Elixir Umbrella应用Oban跨Repo缓存同步问题解决

核心问题定位

你当前的问题根源在于Oban跨应用通知不互通:
Catalog和Orders的Oban实例分别绑定了独立的Repo(Catalog.Repo/Orders.Repo),而Oban.Notifiers.PG依赖PostgreSQL的LISTEN/NOTIFY机制,该机制基于数据库连接池实现——不同Repo的连接池完全独立,导致Orders的监听器根本接收不到Catalog发送的通知,自然无法触发缓存更新。

解决方案

方案1:统一Oban通知器的数据库连接(适用于同数据库场景)

如果两个Repo连接的是同一个数据库,修改Orders的Oban配置,让其通知器复用Catalog的Repo连接,从而共享PG通知通道:

# config/config.exs 或 orders应用的配置文件
config :orders, Orders.Oban,
  repo: Orders.Repo,
  plugins: [
    {Oban.Plugins.Pruner, max_age: 300}
  ],
  # 指定用Catalog的Repo来监听PG通知
  notifier: {Oban.Notifiers.PG, repo: Catalog.Repo},
  queues: [default: 10]

方案2:改用Phoenix PubSub跨应用通信(Umbrella推荐方案)

Umbrella应用内跨应用通信更适合用Phoenix PubSub,它专为Erlang集群/应用间消息传递设计,比Oban PG通知更直接:

  1. 统一PubSub实例
    修改两个应用的Application配置,使用全局共享的PubSub实例:

    # Catalog.Application
    {Phoenix.PubSub, name: MyApp.PubSub}
    
    # Orders.Application
    {Phoenix.PubSub, name: MyApp.PubSub}
    
  2. 替换Catalog的通知发送逻辑
    修改Catalog.Jobs.ProductSyncJob,用Phoenix PubSub发送消息:

    defmodule Catalog.Jobs.ProductSyncJob do
      # ... 原有代码不变
    
      @impl Oban.Worker
      def perform(%Oban.Job{args: %{"product_id" => product_id, "action" => action}}) do
        # ... 原有逻辑不变,替换最后一行的Oban.Notifier.notify为:
        Phoenix.PubSub.broadcast(MyApp.PubSub, "catalog_product_updated", {:product_updated, payload})
        :ok
      end
    end
    
  3. 修改Orders的监听器
    更新Orders.Sync.Listener,监听Phoenix PubSub的消息:

    defmodule Orders.Sync.Listener do
      use GenServer
      require Logger
      alias Phoenix.PubSub
    
      def start_link(_) do
        GenServer.start_link(__MODULE__, nil, name: __MODULE__)
      end
    
      @impl true
      def init(_) do
        Logger.info("Orders.Sync.Listener initializing...")
        PubSub.subscribe(MyApp.PubSub, "catalog_product_updated")
        Logger.info("Listening for catalog_product_updated events")
        {:ok, %{}}
      end
    
      @impl true
      def handle_info({:product_updated, payload}, state) do
        Logger.info("Received product update notification: #{inspect(payload)}")
    
        %{
          product_id: payload.product_id,
          action: payload.action,
          data: payload
        }
        |> Orders.Sync.Jobs.SyncProductJob.new()
        |> Orders.Oban.insert()
    
        {:noreply, state}
      end
    
      @impl true
      def handle_info(msg, state) do
        Logger.debug("Unhandled message in Orders.Sync.Listener: #{inspect(msg)}")
        {:noreply, state}
      end
    end
    

额外排查验证步骤

  1. 检查通知接收日志:查看Orders应用日志,确认是否有Received product update notification的输出——如果没有,说明通知通道仍有问题。
  2. 验证缓存操作日志:给ProductCache的增删改方法添加详细日志,确认操作是否真的执行:
    # 在Orders.ProductCache模块中
    def update_product_cache(product_cache, attrs) do
      Logger.debug("Updating cache for product #{product_cache.catalog_id} with attrs: #{inspect(attrs)}")
      result = product_cache |> ProductCache.changeset(attrs) |> Repo.update()
      Logger.debug("Cache update result: #{inspect(result)}")
      result
    end
    
  3. 确认Oban任务执行:检查Orders的Oban任务日志,确认SyncProductJob是否真的被调度执行。

内容的提问来源于stack exchange,提问作者IvonaK

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 20:24:55