跳到主要内容

职位发布

持久化并检索一条职位发布信息,同时把每个已接受的更新按顺序发送到 LinkedIn 和 Indeed。

产品需求

这个流程在职位发布信息的整个生命周期内代表同一个实体。它必须:

  1. 使用职位名称、描述、最后更新时间和初始值为 0 的更新版本创建 Flow。名称和描述是支持全文检索的 Attribute。
  2. 只运行一次 InitStep,并行启动 UpdateLinkedInPostingUpdateIndeedPosting。之后每个 Step 在自己的 Channel 上等待,且不会占用 Worker。
  3. 通过读取 RPC 返回当前名称、描述和备注。
  4. 通过更新 RPC 接收新的名称、描述和备注。RPC 使用专用且值为 null 的 UpdatePostingLock Attribute,确保同一时间只有一个已接受的更新可以分配版本并发布消息。
  5. 递增 UpdateVersion,持久化当前职位信息,并把同一个不可变 PostingUpdate 发布到 LinkedInPostingUpdatesIndeedPostingUpdates。消息包含版本、稳定的幂等键和职位快照。
  6. 从每个 FIFO Channel 一次消费一条消息。外部调用成功后,consumer 回到自身的新 execution,等待下一条消息。
  7. 独立重试两个招聘网站的调用:首次间隔三秒,最大间隔 60 秒,最多 100 次,总时长不超过一小时。
  8. 删除职位时停止 Flow。搜索操作通过建立索引的 Attribute 查找活跃职位。

当任一 Channel 为空时,Flow 会持久化等待;Step 失败时还会在重试之间等待。更新 RPC 不等待 LinkedIn 或 Indeed。当前状态和两条 Channel 消息持久化提交后,它会返回已提交的版本。

潜在失败包括:

  • 创建时可能缺少必需的初始 Attribute,或者值无法编码。此时不会留下部分创建的 Flow。
  • Flow ID 不存在、Flow 已停止、输入无法持久化,或者另一个 RPC 持有 UpdatePostingLock 时,更新 RPC 会失败。lock conflict 可以安全重试。
  • LinkedIn 或 Indeed 可能不可用、超时、限流或拒绝职位信息。一个 consumer 可以继续运行,而另一个仍在重试或最终失败。
  • 远端请求可能已经成功,但响应丢失。由版本生成的幂等键让外部适配器可以对重试去重。
  • 某个 consumer 失败后,该目标的后续消息会等待它成功或耗尽重试。这可以保持 FIFO 顺序,避免较新的职位状态越过旧状态。
  • 删除可能与已接受的更新竞争。停止策略必须按照产品的删除约定取消或完成排队中和执行中的工作。

lock 决定成功提交的 RPC 顺序,而不是客户端请求到达顺序。并发调用方可能竞争 lock。返回的版本是权威的接受顺序,每个目标都会按照 FIFO 顺序观察这些版本。

Flow 设计

这个交互式 definition graph 由可运行的 Python Flow 生成。

Definition graph

JobPostingFlow

Valid

python · examples/python/dex_examples/products/job-post/job_post_flow.py

这个 Flow 使用一个带 lock 的 producer 和两个持久化、独立推进的 consumer:

创建职位(版本 0)


InitStep
┌────┴────┐
▼ ▼
UpdateLinkedIn UpdateIndeed
Posting Posting
│ │
等待 LinkedIn 等待 Indeed
Channel Channel
▲ ▲
│ │
└────┬────┘

更新 RPC [UpdatePostingLock]
持久化版本 N + 发布相同的 PostingUpdate

└── 返回版本 N

设计使用了以下模式:

  • 实体存储为每条职位信息分配稳定的 Flow ID,并使用 Attribute 保存可检索的当前状态。
  • 初始化扇出使用 InitStep 恰好创建一次两个长期运行的 consumer 分支。
  • 带 lock 的 producer通过 UpdatePostingLock 串行处理版本分配、状态替换和消息发布。
  • Flow 内的事务型 outbox把当前状态和两条 Channel 消息作为一个 RPC 结果提交,因此已接受的更新不会只发给一个目标。
  • 不可变版本消息让每个 consumer 获得触发它的 RPC 中的准确职位快照,而不是稍后读取可变的共享 Attribute。
  • FIFO Channel consumer 循环每次处理一条消息,再回到同一个 Step。即使 RPC 快速连续到达,每个目标也会先观察 v1,再观察 v2。
  • 并行故障隔离让 LinkedIn 和 Indeed 拥有独立的等待、execution、重试历史和进度。较慢的目标不会阻塞另一个目标。
  • 持久化重试处理临时外部故障,在进程重启后也不会丢失进度。
  • 索引状态直接支持名称和描述搜索,不需要应用维护第二份目录。

目标 Step 不需要 Attribute lock。每个目标只有一条串行 consumer 链,顺序由 Channel 提供。增加目标 lock 可以避免重叠,但它本身无法保存更新快照或保证 FIFO 交付。

核心实现

下面的标签页包含每种语言的 Flow 签名和完整机制:InitStep、RPC lock、带版本的双路发布、单消息 Channel 等待、consumer 循环和重试策略。所有片段都来自链接的可运行示例。

class JobPostingFlow(Flow[None]):
update_version = UPDATE_VERSION
update_posting_lock = UPDATE_POSTING_LOCK
linkedin_posting_updates = LINKEDIN_POSTING_UPDATES
indeed_posting_updates = INDEED_POSTING_UPDATES

def get_steps(self) -> StepList[None]:
return StepList.start_step(self.init).other_steps(
self.update_linkedin_posting,
self.update_indeed_posting,
)

@rpc(lock_attributes=(update_posting_lock.lock(),))
def update(self, context: Context, input: JobInfo) -> RPCResult[int]:
version = self.update_version.get(context) + 1
self.title.set(context, input.title or "")
self.job_description.set(context, input.description or "")
self.update_version.set(context, version)
update = PostingUpdate(version, f"{context.flow_id}:{version}", input)
self.linkedin_posting_updates.publish(context, update)
self.indeed_posting_updates.publish(context, update)
return RPCResult(version)


class InitStep(Step[None]):
def execute(self, context: Context, input: None) -> StepDecision:
return go_to_many(
StepMovement.of(UpdateLinkedInPosting, None),
StepMovement.of(UpdateIndeedPosting, None),
)


class UpdateLinkedInPosting(Step[None]):
def wait_for(self, context: Context, input: None) -> Wait:
return Wait.until(LINKEDIN_POSTING_UPDATES.for_one())

def execute(self, context: Context, input: None) -> StepDecision:
update = LINKEDIN_POSTING_UPDATES.results(context)[0]
self.service.update_external_system(
f"update LinkedIn job posting v{update.version} "
f"[{update.idempotency_key}]: {update.posting.title}"
)
return go_to(UpdateLinkedInPosting, None)

例子: examples/python/dex_examples/products/job-post/job_post_flow.py

演示

运行 Dex 和任意语言的示例服务,然后在示例 Playground中创建并更新职位。下面的运行接受了两次更新。执行图显示 InitStep、每个招聘网站的两个有序 execution,以及等待下一条 Channel 消息的 LinkedIn 和 Indeed consumer execution。

职位发布 Flow,其中 InitStep 为每个招聘网站启动两个有序的 Channel consumer execution

相关内容