十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Automatisch 集成开发指南:从零实现一个轮询型 Trigger(以 The cat API 为例)

Automatisch 集成开发指南:从零实现一个轮询型 Trigger(以 The cat API 为例) Automatisch 集成开发指南从零实现一个轮询型 Trigger以 The cat API 为例【免费下载链接】automatischThe open source Zapier alternative. Build workflow automation without spending time and money.项目地址: https://gitcode.com/GitHub_Trending/au/automatisch导读本篇技术指南面向希望为 Automatisch 贡献新应用集成或自定义触发器的开发者以搜索猫图片The cat API为完整实战案例一步步讲解Trigger触发器的目录结构、元数据定义、run函数实现、分页轮询、数据去重internalId以及Test Continue测试流程。阅读本文后你将能够为任意第三方 API 编写一个可被 Automatisch 流程调度器按时轮询、可去重、可正常测试的轮询型触发器并理解其底层执行原理。本文是 Automatisch 官方文档 Build integrations 系列的第 5 篇Triggers建议按顺序阅读以获得完整上下文Folder structure目录结构App应用定义Global variable全局变量Auth认证Triggers触发器即本文Actions动作Examples示例本文全程使用轮询型 HTTP Trigger作为示例。如果你需要实现的是webhook 型触发器例如接收 GitHub 的 webhook 推送请参考 Examples 页面的 Webhook-based triggers 章节。一、把触发器挂载到 App 上Automatisch 中每个应用App由defineApp声明应用下通过auth、triggers、actions三个字段聚合能力。要为thecatapi应用添加触发器先打开thecatapi/index.js加入高亮的import与triggers字段import defineApp from ../../helpers/define-app.js; import auth from ./auth/index.js; import triggers from ./triggers/index.js; export default defineApp({ name: The cat API, key: thecatapi, iconUrl: {BASE_URL}/apps/thecatapi/assets/favicon.svg, authDocUrl: {DOCS_URL}/apps/thecatapi/connection, supportsConnections: true, baseUrl: https://thecatapi.com, apiBaseUrl: https://api.thecatapi.com, primaryColor: #000000, auth, triggers });关于defineApp本身它的实现非常轻量——packages/backend/src/helpers/define-app.js 仅仅是透传应用定义对象export default function defineApp(appDefinition) { return appDefinition; }也就是说应用定义本身是一份纯数据契约key是应用在 Automatisch 内部的唯一标识apiBaseUrl会作为$.http客户端的基础地址这一点在后面的run实现中会用到。二、创建triggers/index.js统一导出在thecatapi目录下创建triggers/index.js将应用内所有触发器以数组形式导出import searchCatImages from ./search-cat-images/index.js; export default [searchCatImages];触发器列表顺序 界面展示顺序。如果你后续新增触发器必须把新触发器追加到这个数组里。Automatisch 前端在选择触发事件步骤中会严格按照数组元素的先后顺序渲染选项。三、定义触发器元数据defineTrigger接下来创建triggers/search-cat-images/index.js用defineTrigger声明触发器元数据与执行逻辑import defineTrigger from ../../../../helpers/define-trigger.js; export default defineTrigger({ name: Search cat images, key: searchCatImages, pollInterval: 15, description: Triggers when there is a new cat image., async run($) { // TODO: Implement trigger! }, });各字段含义如下字段说明name触发器的展示名称显示在 Automatisch 用户界面中。key触发器的唯一标识用于在 Automatisch 内部识别该触发器应保持稳定避免发布后变更。pollInterval轮询间隔分钟。注意该字段目前仅作为声明Automatisch 并未实际读取它——当前轮询间隔默认固定为 15 分钟且不可修改。description触发器的功能描述。run触发器被触发轮询到时间点时执行的核心函数接收全局变量$。3.1defineTrigger的运行时校验defineTrigger并不是简单透传它带有类型校验逻辑。packages/backend/src/helpers/define-trigger.js 的源码如下export default function defineTrigger(triggerDefinition) { const isWebhookOrPoll triggerDefinition.pollInterval || triggerDefinition.type webhook; const schedulerTriggers [ everyNMinutes, everyHour, everyDay, everyWeek, everyMonth, ]; const isSchedulerTrigger schedulerTriggers.includes(triggerDefinition.key); const isMcpTrigger triggerDefinition.key mcpTool; const haveValidTriggerType isWebhookOrPoll || isSchedulerTrigger || isMcpTrigger; if (!haveValidTriggerType) { logger.info(triggerDefinition); throw new Error( Trigger must have a poll interval or be a webhook for ${triggerDefinition.key} ); } return triggerDefinition; }由此可以确认 Automatisch 支持的触发器类型共三类轮询型Polling定义了pollInterval字段如本例的searchCatImages。Webhook 型type: webhook由外部系统推送触发。调度器型Schedulerkey属于everyNMinutes/everyHour/everyDay/everyWeek/everyMonth之一这类触发器用于每 N 分钟/每小时/每天…这类时间驱动场景仓库中的实现见 packages/backend/src/apps/scheduler/triggers/every-n-minutes/index.js它不轮询第三方 API而是通过getInterval(parameters)返回 cron 表达式驱动调度。如果触发器既没有pollInterval、又不是 webhook、也不属于调度器/MCP 类型defineTrigger会直接抛出错误避免把无效配置带入运行时。四、实现run函数轮询、分页与数据推送现在补全run的实现这是整个触发器开发的核心import defineTrigger from ../../../../helpers/define-trigger.js; export default defineTrigger({ // ... async run($) { let page 0; let response; const headers { x-api-key: $.auth.data.apiKey, }; do { let requestPath /v1/images/search?page${page}limit10orderDESC; response await $.http.get(requestPath, { headers }); response.data.forEach((image) { const dataItem { raw: image, meta: { internalId: image.id }, }; $.pushTriggerItem(dataItem); }); page 1; } while (response.data.length 10); }, });4.1$.http内置 HTTP 客户端代码中的$.http来自 Automatisch 的全局变量机制。在 packages/backend/src/engine/global-variable.js 中它以应用定义里的apiBaseUrl为baseURL创建客户端并自动挂载应用的beforeRequest钩子如 OAuth 签名、Token 注入等。因此本例中请求路径只需写相对路径/v1/images/search...且apiBaseUrl已在上文defineApp中配置为https://api.thecatapi.com。$.auth.data则来自用户建立连接Connection时保存的认证数据apiKey这正是前面 Auth 篇 的产出物。4.2 分页轮询的两种常见写法本例采用do...while循环从第 0 页开始每页取 10 条只要本页仍返回满 10 条数据就继续下一页直到某页不足 10 条即视为到达最后一页do { // 请求第 page 页limit10 response await $.http.get(requestPath, { headers }); // 逐条推送 page 1; } while (response.data.length 10);仓库中的真实集成则展示了另一种基于Link 响应头的分页方式。例如 GitHub 的 new-issues 触发器通过解析响应头里的link字段判断是否存在next页let links; do { const response await $.http.get(pathname, { params }); links parseLinkHeader(response.headers.link); if (response.data.length) { for (const issue of response.data) { const dataItem { raw: issue, meta: { internalId: issue.id.toString() }, }; $.pushTriggerItem(dataItem); } } } while (links.next);两种写法本质相同在单次run中尽可能把当前能拿到的新数据全部抓取完并逐条推送由 Automatisch 负责去重见下一节。4.3 两条铁律反向时间序 internalIdrun函数本身不需要返回值数据完全依靠$.pushTriggerItem推送给 Automatisch。但在编写时有两条规定必须严格遵守否则触发器无法正常工作按倒序reverse-chronological order推送数据。轮询型触发器的去重机制依赖最早遇到已处理数据就停止这一假设因此必须把最新数据排在最前面否则可能漏掉中间的新数据。必须为每条数据提供唯一标识internalId。internalId可以是第三方数据的 ID也可以是任何能唯一标识该条数据的字段。Automatisch 用它判断数据是否已被处理从而防止重复数据进入后续步骤。$.pushTriggerItem接受的对象结构为字段说明raw推送给 Automatisch 的原始数据通常是第三方 API 返回的原始对象。meta数据元信息必须包含internalId字段。五、$.pushTriggerItem的底层原理理解$.pushTriggerItem的行为能帮你写出更健壮的触发器。其实现位于 packages/backend/src/engine/global-variable.jspushTriggerItem: (triggerItem) { if ( isAlreadyProcessed(triggerItem.meta.internalId) !$.execution.testRun ) { // early exit as we do not want to process duplicate items in actual executions throw new AlreadyProcessedError(); } $.triggerOutput.data.push(triggerItem); const isWebhookApp app.key webhook; const isFormsApp app.key forms; if ($.execution.testRun !isWebhookApp !isFormsApp) { // early exit after receiving one item as it is enough for test execution throw new EarlyExitError(); } },结合这一实现可以得出几个重要的行为事实5.1 已处理数据会自动中断执行全局变量在构建时会加载该 flow 的历史internalId集合const lastInternalIds testRun || (flow step?.isAction) ? [] : await flow?.lastInternalIds(2000); const isAlreadyProcessed (internalId) { return lastInternalIds?.includes(internalId); };flow.lastInternalIds从该 flow 最近的 executions 中取出最多 2000 条internal_id见 packages/backend/src/models/flow.js。当$.pushTriggerItem收到的internalId命中这个集合、且当前不是测试运行时会抛出AlreadyProcessedError从而停止触发器——而此前已经通过$.triggerOutput.data.push(...)推入的数据会照常进入后续处理。去重的作用域是按 flow 隔离的同一个第三方 ID 在不同流程中互不影响。5.2 测试运行Test Continue的特殊行为$.pushTriggerItem能自动区分用户点击Test Continue按钮触发的测试执行与已发布流程的正式调度执行测试执行推入第一条数据后立即抛出EarlyExitError提前退出无论该数据此前是否已被处理过测试场景下lastInternalIds为空不做去重判断。正式执行会持续抓取并推送剩余数据直到遇到已处理数据或分页结束。唯一的例外是webhook与forms两个应用——它们的触发器在测试运行时不提前退出。仓库中 webhook 触发器的实现可见 packages/backend/src/apps/webhook/triggers/catch-raw-webhook/index.js它还额外实现了testRun函数测试时直接复用上一次执行成功保留的dataOut而不是等待真实的外部请求。5.3 中途失败不会丢失已抓取的数据一个典型的容错场景假设触发器已经成功抓取了前 5 页数据发了 5 个 HTTP 请求在请求第 6 页时遇到第三方 API 的限流rate limit错误。此时 Automatisch不会丢失前 5 页已经抓取到的数据——它会在第一次遇到错误时停止触发器但会把所有已经成功推入的数据照常处理完。这是边抓取边推送 数据先入triggerOutput.data再统一处理这一设计带来的天然容错性。5.4 internalId 的落库每条被执行的数据的internalId会写入对应 execution 记录。见 packages/backend/src/engine/trigger/process.jsconst execution await Execution.query().insert({ flowId, testRun, internalId: initialDataItem?.meta.internalId, });这也解释了为什么按倒序推送是硬性要求Automatisch 需要靠internalId判定该从哪条数据之后开始算新数据顺序错误会导致新旧边界判断失误。六、调度器如何驱动轮询型触发器轮询型触发器由 Automatisch 的调度队列按 cron 表达式周期驱动。packages/backend/src/models/flow.js 中的getExecutionIntervalAsCron展示了间隔到 cron 的映射const EVERY_1_MINUTE_CRON * * * * *; const EVERY_2_MINUTES_CRON */2 * * * *; const EVERY_5_MINUTES_CRON */5 * * * *; // ... 10 / 15 / 30 / 60 分钟对应 cron const EVERY_60_MINUTES_CRON 0 * * * *; async getExecutionIntervalAsCron() { const interval this.executionInterval || 15; if (!(await hasValidLicense())) { return EVERY_15_MINUTES_CRON; // 无有效 License 时固定 15 分钟 } switch (interval) { case 1: return EVERY_1_MINUTE_CRON; // ... case 60: return EVERY_60_MINUTES_CRON; default: return EVERY_15_MINUTES_CRON; } }同时 flow 模型的 jsonSchema 将executionInterval限定为[1, 2, 5, 10, 15, 30, 60]默认15。这从源码层面印证了官方文档中的说明当前轮询间隔默认固定为 15 分钟且只有在持有有效 License 时才能把间隔调整为更细粒度1/2/5/10 分钟触发器定义里的pollInterval字段目前确实只是声明性元数据。当流程发布updateStatus将active置为true时轮询型触发器通过Engine.runInBackground注册带repeat选项的定时任务webhook 型触发器则调用registerHook注册远端 webhook见 packages/backend/src/models/flow.js。若需要 webhook 型触发器的完整开发细节务必阅读 Examples 页面的 Webhook-based triggers 章节。七、测试触发器Test Continue实现完成后进入 Automatisch 界面验证打开Flows流程页面新建一个流程选择The cat API应用再选择Search cat images触发器点击Test Continue按钮。如果界面上能看到 JSON 响应数据说明触发器工作正常。测试阶段$.pushTriggerItem只推送第一条数据即提前退出见 5.2 节这既保证了响应速度也足够你验证请求路径、认证头和数据结构是否正确。八、开发核对清单检查项要求挂载在thecatapi/index.js的defineApp中加入triggers字段导出triggers/index.js以数组导出全部触发器顺序即界面顺序类型defineTrigger定义中必须含pollInterval或type: webhook调度器/MCP 除外否则启动报错顺序run内按倒序最新在前推送数据去重每条数据meta.internalId必须唯一且稳定分页建议一次run抓完当前所有新数据do...while或解析 Link 头容错抓取中途失败不会丢已抓数据无需手动缓存分页游标测试用Test Continue验证首条数据能否正确返回按上述清单完成实现后你的轮询型触发器就可以随应用一起交付使用。更完整的真实集成实现含 webhook 型与调度器型可继续阅读本系列的 Examples 篇或直接在仓库packages/backend/src/apps/下查看各应用triggers/目录中超过 100 个已实现触发器作为参照。【免费下载链接】automatischThe open source Zapier alternative. Build workflow automation without spending time and money.项目地址: https://gitcode.com/GitHub_Trending/au/automatisch创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表