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通知更直接:
统一PubSub实例
修改两个应用的Application配置,使用全局共享的PubSub实例:# Catalog.Application {Phoenix.PubSub, name: MyApp.PubSub} # Orders.Application {Phoenix.PubSub, name: MyApp.PubSub}替换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修改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
额外排查验证步骤
- 检查通知接收日志:查看Orders应用日志,确认是否有
Received product update notification的输出——如果没有,说明通知通道仍有问题。 - 验证缓存操作日志:给
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 - 确认Oban任务执行:检查Orders的Oban任务日志,确认
SyncProductJob是否真的被调度执行。
内容的提问来源于stack exchange,提问作者IvonaK
相关产品推荐
相关产品推荐

