IT加油站

爬虫不等于监控:用 Apify 和 n8n 检测 Telegram 新帖

11浏览 3小时前 软件教程 MA123123

原文:https://dev.to/apify/a-scraper-isnt-a-monitor-detecting-new-telegram-posts-with-apify-and-n8n-31f5(作者 @bhagyarathnasekara)

我的 Actor 运行成功了。它返回了五条最近的 Telegram 帖子,每条都带 ID、时间戳和永久链接。

于是我提出了这个自动化真正需要回答的问题:

哪些帖子是新的?

诚实的答案是:当前这次运行无法自行判断。

没有任何环节失败。Actor 完成了它的工作:它取回了当前可见的内容。但“新”并不是单个快照所包含的属性;它是当前快照与之前某个快照之间的差值。

这个区别,正是“按计划运行的爬虫”和“可以信赖的监控器”之间的差异。

在本教程中,我会展示我为 Apify Actor 添加的状态层:

  • 每条记录的稳定身份标识;
  • 按 Telegram 频道划分的状态;
  • 审慎设计的首次运行策略;
  • 跨执行去重;
  • 当轮询窗口可能过小时给出警告;
  • 对缺失或损坏状态采取 fail-closed(默认安全)处理。

同样的设计也适用于价格追踪器、职位监控、销售线索推送、库存检查,以及其他任何反复追问“发生了什么变化?”的工作流。

实现说明: 本文的代码和验证材料位于 修正版配套包。该工作流有意止步于“通知候选项”;我没有把未经测试的投递当作已完成。

获取与变更检测是两件不同的事

该 Actor 读取 Telegram 公开的 t.me/s/ 预览,并按最新优先的顺序返回近期消息。每条消息包含如下字段:

{
  "channel": "telegram",
  "id": 454,
  "url": "https://t.me/telegram/454",
  "date": "2026-07-19T17:58:20+00:00",
  "text": "For all the details on these new features...",
  "views": 1180000,
  "scraped_at": "2026-08-05T14:13:17.444080+00:00"
}


这足以回答:

现在能看到什么?

却不足以回答:

自上次成功检查以来出现了什么?

第二个问题需要历史信息。这段历史可以放在 Actor 内部、数据库里、调用方提供的游标中,也可以放在调用 Actor 的工作流里。在本次实现中,我让 Actor 保持无状态,把比较状态交由下游 n8n 工作流负责。

边界看起来像这样:

定时调度
   │
   ▼
运行 Apify Actor ──► 当前消息窗口
                         │
                         ▼
             校验并分区记录
                         │
                         ▼
             与持久化已见状态比较
                         │
              ┌──────────┴──────────┐
              ▼                     ▼
         通知候选项              状态与警告
              │                     │
              └──────────┬──────────┘
                         ▼
                 持久化新状态


暴露缺失层的实验

我最初通过 Apify MCP server 从 AI 客户端调用 Actor,并保留了四种“单次调用”场景:

| 场景 | 运行前调用方知道什么 | 能否判断新 ID? |

| --- | --- | --- |

| 基线(Baseline) | 没有之前的结果 | 未要求判断;本次结果成为基线 |

| 同一会话 | 之前的 ID 仍在上下文中 | 能 |

| 全新会话 | 没有可访问的先前结果 | 不能;新属性无法判定 |

| 带先前 ID 的全新会话 | 一小份状态文件 + 比较规则 | 能 |

四次 Actor 运行全部成功,返回的都是相同的五个 ID,从 454450。在没有基线的全新会话中,获取仍然成功,但比较没有完成。

这个结果很重要,因为 “没有新记录”和“我无法判断记录是否是新记录”并不是同一个答案。 可靠的监控器必须保留这个区分。

这几次运行还暴露了一个很有迷惑性的身份标识 bug:两次观察之间,每条消息的 views 值都变了,而 ID 保持稳定。如果比较整条记录,五条消息都会被判定为“不同”;如果比较稳定身份标识,则能正确判定它们是同一些消息。

先选身份标识,再选存储

Telegram 的 Bot API 将 message_id 定义为聊天内的唯一标识。也就是说,单独的 id 不能作为可复用的全局键。因此监控器改用这种分区身份:

(normalized channel, message ID)


例如:

function normalizeChannel(value) {
  return String(value ?? '')
    .trim()
    .replace(/^@/, '')
    .toLowerCase();
}

function messageKey(channel, id) {
  const normalizedChannel = normalizeChannel(channel);
  if (!normalizedChannel) {
    throw new Error('Channel is required.');
  }

  if (!Number.isSafeInteger(id) || id <= 0) {
    throw new Error('Message ID must be a positive safe integer.');
  }

  return JSON.stringify([normalizedChannel, String(id)]);
}


viewsscraped_at 被有意排除在外。它们可以变化,而不代表产生了新消息。文本也可能被编辑,所以仅基于 ID 的状态能检测新消息,但检测不到编辑;编辑检测是另一份独立契约。

这是第一条设计规则:

用不可变的源身份标识来检测新条目。不要让可变属性重新定义记录。

存储基线,而非模糊记忆

单个通道的最小先前状态文件可以小到:

{
  "schemaVersion": 1,
  "channel": "telegram",
  "identityField": "id",
  "sourceRunId": "example-prior-run-id",
  "seenMessageIds": [454, 453, 452, 451, 450]
}


n8n 工作流以更紧凑的形式使用同一契约。其全局工作流状态按规范化通道分区,正整数 ID 以合并区间存储:

{
  "schemaVersion": 1,
  "channels": {
    "telegram": {
      "initializedAt": "<timestamp from the first successful run>",
      "baselineFloorId": 450,
      "seenRanges": [[450, 454]],
      "highestObservedId": 454
    }
  }
}


区间压缩是实现层面的优化,不属于通用监控规则的一部分。集合比较需要稳定标识符;区间压缩额外要求 ID 是合适的整数。

工作流通过 n8n 的 $getWorkflowStaticData('global') 读取该状态。n8n 文档说明了三个重要约束:静态数据应保持较小规模,在成功的触发式生产执行后保存,并且在高频执行下可能不可靠。因此它只适合这个小型序列化演示,不能当作通用生产数据库。

对于持续增长的状态、重叠执行或运维查询场景,应将同样的标识与比较契约迁移到 Data Table 或外部数据库。

将首次运行视为策略决策

首次成功执行时,返回的每个 ID 都是未见过的。把所有 ID 都称作“新”在数学上一致,但通常是糟糕的监控行为:工作流会对监控开始前就已存在的帖子发送一波告警。

该实现默认采用基线模式:

  1. 获取第一个非空窗口。
  2. 校验它。
  3. 存储这些 ID。
  4. 不发送任何通知候选项。
  5. 在下次成功运行时开始新条目检测。

对于确实希望交付首个窗口的工作流,仍可显式使用 alert 模式。

空结果不会初始化基线。否则一次瞬时的空响应会创建空历史,下一次正常结果就会把整个窗口当作新条目重放。

比较、分类,再持久化

核心比较是集合差:

新 ID = 当前 ID − 已见 ID


但该表达式周围的顺序至关重要:

  1. 校验响应结构和通道分区。
  2. 拒绝不可用的 ID。
  3. 按分区身份对当前窗口去重。
  4. 与已见身份进行比较。
  5. 应用首次运行和历史回填策略。
  6. 生成通知候选项。
  7. 持久化每个有效的已观察 ID,包括被关键词过滤掉的记录。

最后这一选择让关键词过滤器具有最多一次(at-most-once)行为。之后修改关键词不会重放工作流已处理过的帖子。

状态损坏时也会安全失败(fail closed)。如果非空状态对象没有受支持的 schema,或包含格式错误的通道状态,工作流会停止,而不是默默把历史视为空。重置损坏的状态会把旧记录变成新一轮告警风暴。

完整的比较器与 n8n Code 节点实现位于 公开配套仓库。比较器无依赖,既接受裸数据集数组,也接受 MCP 风格的 { "items": [...] } 封装结构。

将轮询窗口满载视为警告

假设工作流请求五条消息并收到五条。它无法确定是恰好有五条可用,还是有更早的未读消息落在返回窗口之外。

因此工作流使用 Actor 原始响应计数设置 windowSaturated

const windowSaturated = rows.length >= Number(monitor.limit);


该检查发生在去重之前。重复或无效的行不应掩盖原始窗口已达上限的事实。

窗口满载意味着:

配置的窗口对于轮询间隔来说可能太小。

并不证明有消息被漏掉。它提示运维人员考虑更大的限制、更短的间隔或合适的追赶策略。

同样的推理也适用于相反方向:一个此前已见、但当前满载窗口中没有出现的 ID 不会被自动删除。它可能只是落到了窗口之外。

MCP 验证了信息边界,n8n 负责周期状态

本项目中的两条调用路径有意分离:

  • 实验使用 AI 客户端和 Apify MCP 服务器,测试在有/无先前状态时分别能确定什么。
  • 周期实现使用 n8n 和 Apify 的同步 REST 端点获取数据集条目,与 n8n 工作流状态比较,并生成通知候选项。

它们调用同一个 Actor,但并不是一条 MCP 到 n8n 的组合流水线。

在 n8n 中,HTTP Request 节点调用:

POST https://api.apify.com/v2/acts/{ACTOR_ID}/run-sync-get-dataset-items


Apify API 将该端点文档描述为同步运行 Actor 并返回其数据集条目。请将 Apify API token 存放在 n8n 的 Bearer Auth 凭据中,不要直接粘贴到工作流 JSON 里。

我验证了什么

修正后的配套包包含:

  • 一个零依赖的比较器,附带23 项确定性测试
  • 一个可导入的 n8n 工作流,其内嵌状态引擎附带11 项确定性测试
  • 故障关闭(fail-closed)的状态模式;
  • 首次运行基线抑制;
  • 按渠道划分的状态分区;
  • 跨触发运行的重复预防;
  • 基于原始响应数量的饱和报告。

我还保留了两次成功的 n8n 触发执行记录

| 执行 | 结果 |

| --- | --- |

| 第一次触发运行 | 为两个渠道创建了独立基线;零个通知候选项 |

| 第二次触发运行 | 复用了持久化状态;返回的十个 ID 全部已见 |

第二次执行读取第一次执行的状态,这是最重要的结果。它证明了持久化在触发的生产运行之间有效,而不仅仅是在某个编辑器会话内部有效。

我没有验证什么

证据存在边界:

  • 在保留的运行期间,没有出现真正新的实时 Telegram ID。新 ID 的正向路径由确定性构造测试覆盖。
  • 修正后的公开工作流止步于通知候选项;Telegram 消息发送功能被故意禁用且未测试。
  • 每个观察结果窗口都达到了配置的 5 条上限,因此该实验无法证明追赶补全的完整性。
  • n8n 工作流静态数据仅使用小状态和序列化执行进行过测试。
  • MCP 实验与 n8n 工作流在不同时间运行,并非单变量性能基准。

这些边界并不会削弱状态契约。它们恰恰界定了该实现证明了什么——以及下一次生产测试必须覆盖什么。

我现在每次调度 Actor 前都会用的检查清单

在把任何周期性 Actor 调用变成监控器之前,我都会回答这些问题:

  1. 稳定身份是什么? 必要时按来源、租户、渠道或账号进行分区。
  2. 哪些字段是可变的? 除非变更检测本身就是目标,否则将它们排除在身份之外。
  3. 基线存放在哪里? 让某一个层级拥有明确的所有权。
  4. 首次非空运行时会发生什么? 建立基线或发出告警——绝不能是意外行为。
  5. 状态缺失或损坏时会发生什么? 故障关闭,而不是重放历史。
  6. 轮询窗口会饱和吗? 在去重之前先暴露这一状况。
  7. 缺失意味着什么? 在受限窗口中,缺失并不证明删除。
  8. 执行会重叠吗? 如果会,使用具有明确并发保证的存储。

一个定时抓取器只会反复告诉你当前存在什么。一个监控器则告诉你什么发生了变化,同时不会编造确定性、重放旧记录或掩盖缺口。

差异不在于多一次 API 调用,而在于状态契约。

资源

_AI 辅助披露:我在起草、编辑和 QA 过程中使用了 AI 工具。我已审阅最终文章及支撑证据,并对其主张和结论负责。_

原文:https://dev.to/apify/a-scraper-isnt-a-monitor-detecting-new-telegram-posts-with-apify-and-n8n-31f5(作者 @bhagyarathnasekara)

#Apify #n8n #Telegram #爬虫 #自动化工作流 #状态管理