ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

iii 实战教程:用队列与发布订阅让短链服务具备持久化执行能力(Durable Execution)

iii 实战教程:用队列与发布订阅让短链服务具备持久化执行能力(Durable Execution) iii 实战教程用队列与发布订阅让短链服务具备持久化执行能力Durable Execution【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii导读本文是 iii 项目 Linkly 短链服务系列教程的第 4 章Make it durable。前三章中每次重定向都会直接触发link::record_click写数据库导致数据库写入落在重定向的热路径上——一次慢写入就会拖慢整个重定向响应。本章将这一工作迁移到队列上让重定向立即返回随后用**发布/订阅pub/sub**把链接事件广播给互相解耦的独立订阅者一个 Python 分析 worker 和一个缓存刷新器。读完本文你将掌握TriggerAction.Enqueue的用法、queueworker 的队列配置、pubsub与iii::durable::publish两种发布订阅模式的取舍以及如何用 Python 编写一个跨语言、跨 worker 的事件消费者。原文出处docs/tutorials/linkly/durable-execution.mdx本教程的.skill.md渲染版本为 durable-execution.mdx.skill.md。本章目标把写库工作移出热路径回顾上一章Ch. 3的实现http::redirect直接触发link::record_click也就是每次用户访问短链时重定向函数都要同步等待一次数据库写入完成。这带来两个问题延迟耦合数据库写入慢时重定向也跟着变慢用户体验直接受损责任耦合重定向这个读路径上背负了写日志的额外职责。本章的解法是用两个既有的 worker 能力重构队列queue接受现在提交、稍后执行的工作让link::record_click在后台排空drain队列重定向立刻返回 302发布订阅pub/sub队列把每条消息投递给一个消费者而当系统里多个互不相关的部分需要对同一事件做出反应时改用发布/订阅设计把事件广播给所有订阅者。最终得到两个完全解耦的消费者一个Python analytics worker按天统计短链创建数和一个缓存刷新器link::on_link_updatedlinkworker 甚至不需要知道它们的存在。添加 worker加入 pubsub本章涉及两个 worker。queue在iii project init初始化项目时就已经内置无需重复添加pubsub则需要显式加入。在项目目录下执行iii worker add pubsub这条命令会把pubsubworker 注册到当前项目。注意queue自第 1 章起就在运行它的设置统一管理在./config/queue.yaml中详见 Configuration。用队列让重定向变快队列的基本概念队列保存现在接受、以后运行的工作。本章定义了一个名为clicks的队列专门用于link::record_click。定义 clicks 队列编辑 config/queue.yaml在./config/queue.yaml的value:下添加queue_configs并保存。配置热生效无需重启id: queue # ... value: queue_configs: clicks: type: standard max_retries: 5 concurrency: 5 adapter: name: builtin各字段说明结合 engine/src/workers/queue/README.md 中的配置表字段类型说明默认值typestringstandard并发处理或fifo按消息组有序处理—max_retriesu32进入死信队列DLQ前的最大投递尝试次数3concurrencyu32同时处理的最大任务数FIFO 队列会强制为prefetch110message_group_fieldstring仅fifo队列需要指定决定顺序组的 JSON 字段—backoff_msu64重试基础退避毫秒数指数退避backoff_ms × 2^(attempt−1)1000poll_interval_msu64worker 轮询间隔毫秒100adapter.namestring传输适配器默认builtinbuiltin注意队列名只是对队列的引用不会对放入该队列的内容施加任何限制——你可以把任何函数调用入队到任何已声明的队列。让 record_click 变成可排队任务你在第 3 章已经写好了link::record_click函数本身一行都不用改——可排队性只取决于触发方式。只需用TriggerAction改变它的触发方式。第一步在link/src/index.ts中导入TriggerActionimport { registerWorker, TriggerAction } from iii-sdk; import { Logger } from iii-dev/helpers/observability;第二步给http::redirect中现有的link::record_click调用加上action让queueworker 把它入队而不是内联执行worker.registerFunction(http::redirect, async (req) { // ...previous code... await worker.trigger({ function_id: link::record_click, payload: { code, clicked_at: new Date().toISOString() }, action: TriggerAction.Enqueue({ queue: clicks }), }); return { status_code: 302, headers: { Location: url } }; });从源码看TriggerAction.Enqueue({ queue })会序列化为{ type: enqueue, queue }的线上格式SDK 的契约测试对此有明确锁定见 sdk/packages/node/iii/tests/trigger-action.test.ts引擎的TriggerAction反序列化依赖这个格式。效果改造后重定向在点击被接受进队列的瞬间就返回 302不再等待数据库写入完成。link::record_click在后台排空队列并且自带重试机制如果写入持续失败任务最终进入死信队列DLQ。值得一提的是队列能力不再是引擎内置从源码看引擎会把TriggerAction.Enqueue与durable:subscriber路由到独立的queueworker见 engine/src/workers/queue/README.md如果你遇到enqueue_error: engine::queue::enqueue not found说明项目里没有这个 worker需要把它加进 Compose。用发布订阅广播事件队列把每条消息投递给一个消费者。当系统里多个互不相关的部分需要对同一事件做出反应时应该改用发布/订阅publish/subscribe设计。queue 与 pubsub 的职责划分项目同时提供了queue和pubsub两个 worker官方文档给出了清晰的选择标准需要保证成功或失败进入 DLQ的发布订阅流使用queue提供的iii::durable::publish与durable:subscriber触发类型不需要保证的发布订阅流使用pubsub提供的publish与subscribe。本章会实现两个主题创建链接时发布link.created更新链接时发布link.updated。由于目前还没有链接更新功能我们会一并补上它和一个 HTTP 端点。在 link.created 上发布在link::create内部完成数据库写入和state::set之后触发内置的publish函数worker.registerFunction(link::create, async (payload: { url: string; code?: string }) { // ...previous code... await worker.trigger({ function_id: publish, payload: { topic: link.created, data: { code, url } }, }); logger.info(link created, { code, url }); return { code, url }; });在 link.updated 上发布首先给linkworker 增加更新路径让链接目标可以改变并对外广播。下面是领域函数它更新数据库行并通过持久化发布订阅iii::durable::publish由queueworker 提供发布link.updated事件。把它放在link/src/index.ts末尾worker.registerFunction(link::update, async (payload: { code: string; url: string }) { const url /^https?:\/\//i.test(payload.url) ? payload.url : https://${payload.url}; await worker.trigger({ function_id: database::execute, payload: { db: DB, sql: UPDATE links SET url ? WHERE code ?, params: [url, payload.code], }, }); await worker.trigger({ function_id: iii::durable::publish, payload: { topic: link.updated, data: { code: payload.code, url } }, }); return { code: payload.code, url }; });通过 HTTP 暴露链接更新接着用 HTTP handler 连接link::update负责校验输入并调用领域函数worker.registerFunction(http::update, async (req) { const code req.path_params.code; const url req.body?.url; if (!url) { return { status_code: 400, body: { error: missing url }, headers: { Content-Type: application/json }, }; } const link await worker.trigger{ code: string; url: string }, { code: string; url: string }({ function_id: link::update, payload: { code, url }, }); return { status_code: 200, body: link, headers: { Content-Type: application/json } }; });最后是绑定http::update到PUT /links/:code的触发器worker.registerTrigger({ type: http, function_id: http::update, config: { api_path: /links/:code, http_method: PUT }, });这段代码本身与 pub/sub 无关但它是本章末尾验证新功能所必需的。响应式状态不耦合地保持缓存正确link::update只改了数据库没有改状态缓存因此查询可能读到过期数据。与其在link::update内部处理缓存刷新不如用持久化订阅者订阅link.updated事件worker.registerFunction(link::on_link_updated, async (data: { code: string; url: string }) { await worker.trigger({ function_id: state::set, payload: { scope: links, key: data.code, value: { url: data.url } }, }); }); worker.registerTrigger({ type: durable:subscriber, function_id: link::on_link_updated, config: { topic: link.updated }, });持久化 vs 普通发布订阅这里有一个值得反复琢磨的设计决策link.updated使用持久化发布订阅iii::durable::publish搭配durable:subscriber触发器两者都由queueworker 提供。像缓存刷新器这样的消费者必须收到每一次更新——丢失一个事件缓存就会一直指向过期 URLlink.created留在普通发布订阅pubsub上因为它的唯一消费者是一个尽力而为的每日计数器偶尔漏掉一条无伤大雅。通用原则当错过事件会破坏状态时用持久化发布订阅当只是 fire-and-forget 的扇出时用普通发布订阅。从queueworker 的源码文档看持久化发布订阅本质上是基于主题的队列每个订阅了某主题的函数都会收到每条消息的一份拷贝支持重试与 DLQ见 engine/src/workers/queue/README.md 的 Topic-based queues 说明。用 Python 创建分析 worker队列和事件在单个 worker 内部有用跨 worker 同样有用。此前所有代码都是 TypeScript但 iii 的 worker 不限制语言或运行时——这次我们用 Python 写一个统计链接数量的 analytics worker。创建新 worker用与第 1 章创建linkworker 相同的方式创建 Python worker。它会生成一个analytics/worker包含src/main.py示例和iii.worker.yaml清单iii worker init analytics --language python订阅 link.created 事件把示例src/main.py替换为下面这份代码订阅link.created事件统计每个新短链的创建次数import os from datetime import datetime, timezone from iii import register_worker, InitOptions from iii_helpers.observability import Logger worker register_worker( os.environ.get(III_URL, ws://localhost:49134), InitOptions(worker_nameanalytics), ) logger Logger() DB analytics def ensure_schema() - None: The analytics worker owns its own table, in its own database. worker.trigger( { function_id: database::execute, payload: { db: DB, sql: CREATE TABLE IF NOT EXISTS daily_link_counts (day TEXT PRIMARY KEY, count INTEGER NOT NULL), }, } ) def on_link_created(data: dict) - dict: Runs whenever link publishes link.created. Counts links per day. day datetime.now(timezone.utc).strftime(%Y-%m-%d) worker.trigger( { function_id: database::execute, payload: { db: DB, sql: INSERT INTO daily_link_counts (day, count) VALUES (?, 1) ON CONFLICT(day) DO UPDATE SET count count 1, params: [day], }, } ) logger.info(fcounted new link {data.get(code)} for {day}) return {ok: True} ensure_schema() worker.register_function(analytics::on_link_created, on_link_created) worker.register_trigger( { type: subscribe, function_id: analytics::on_link_created, config: {topic: link.created}, } ) print(Analytics worker started)这段代码值得注意的几点连接建立register_worker默认连接ws://localhost:49134可用III_URL环境变量覆盖Python SDK 中register_worker与InitOptions的实现见 sdk/packages/python/iii/src/iii/iii.py自持数据analytics worker 在自己的数据库、自己的表里记账linkworker 完全不需要知道它的存在幂等建表ensure_schema()在启动时执行CREATE TABLE IF NOT EXISTS按天计数用ON CONFLICT(day) DO UPDATE SET count count 1实现 upsert 式累加。配置 worker现有清单analytics/iii.worker.yaml已经够用name: analytics runtime: # Base OCI image used as the worker rootfs. base_image: docker.io/iiidev/python:latest scripts: install: pip install -e . start: watchfiles python src/main.py字段含义runtime.base_image是 worker 的根文件系统基础 OCI 镜像scripts.install在 worker 运行时内安装 SDK、可观测性辅助库与源码监视器scripts.start用watchfiles启动python src/main.py并支持热重载。analytics 的数据放在自己的数据库里因此需要给databaseworker 的配置文件追加一个analytics数据库与第 3 章的primary并列。编辑config/database.yaml在value: databases:下新增第二项保存后自动生效无需重启id: database name: Database value: databases: primary: pool: acquire_timeout_ms: 5000 idle_timeout_ms: 30000 max: 10 url: sqlite:./data/iii.db analytics: url: sqlite:./data/analytics.db最后把新的 analytics worker 加入项目iii worker add ./analytics验证效果创建五个链接、跟随其中一个几次并修改它的目标# Make some new links for n in $(seq 1 5); do curl -s -X POST http://127.0.0.1:3111/links \ -H Content-Type: application/json -d {\url\:\https://iii.dev/$n\,\code\:\analyticslink$n\} donePython worker 会如预期地记录新建链接的数量。用iii trigger直接查询 analytics 数据库验证iii trigger database::query dbanalytics sqlSELECT day, count FROM daily_link_counts{ rows: [{ day: 2026-05-27, count: 5 }], row_count: 1 }iii trigger也可用于iii::queue::redrive之类的内置函数例如把某命名队列的死信消息重新投回主队列见 engine/src/workers/queue/README.md。适配器选择从本地开发到多实例生产虽然本章的config/queue.yaml只用了默认的builtin适配器但从queueworker 的源码文档engine/src/workers/queue/README.md可以看到适配器是决定队列行为的关键一环值得在选择时留意适配器重试DLQFIFO命名队列消费主题 pub/sub多实例外部依赖builtin是是是是是否无rabbitmq是是是是是是RabbitMQredis否否否否仅发布是是Redis经验法则本地开发用builtinin_memory存储单实例生产用builtinfile_based存储多实例生产用rabbitmq。builtin适配器可配置store_method: in_memory | file_based与file_pathrabbitmq需要amqp_url。多实例场景下若只做主题广播redis适配器也可作为轻量选择。总结本章完成了一次典型的持久化执行重构重定向不再等待数据库写入点击记录乘上clicks队列在后台排空配合重试与死信队列兜底链接事件经 pub/sub 扇出link.created走普通pubsublink.updated走queue提供的持久化发布订阅各自按错过事件是否会破坏状态这一标准选型两个消费者完全解耦Python analytics worker按天计数和缓存刷新器link::on_link_updated都只关心事件主题linkworker 不知道它们的存在跨语言协作同一套触发机制同时服务 TypeScript 与 Python worker印证了 iii worker 不绑定语言或运行时的设计。下一步是第 5 章实时流式点击届时会用专门的click-streamerworker 实时广播每一次点击。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表