原文: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,从 454 到 450。在没有基线的全新会话中,获取仍然成功,但比较没有完成。
这个结果很重要,因为 “没有新记录”和“我无法判断记录是否是新记录”并不是同一个答案。 可靠的监控器必须保留这个区分。
这几次运行还暴露了一个很有迷惑性的身份标识 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)]);
}
views 和 scraped_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 都称作“新”在数学上一致,但通常是糟糕的监控行为:工作流会对监控开始前就已存在的帖子发送一波告警。
该实现默认采用基线模式:
- 获取第一个非空窗口。
- 校验它。
- 存储这些 ID。
- 不发送任何通知候选项。
- 在下次成功运行时开始新条目检测。
对于确实希望交付首个窗口的工作流,仍可显式使用 alert 模式。
空结果不会初始化基线。否则一次瞬时的空响应会创建空历史,下一次正常结果就会把整个窗口当作新条目重放。
比较、分类,再持久化
核心比较是集合差:
新 ID = 当前 ID − 已见 ID
但该表达式周围的顺序至关重要:
- 校验响应结构和通道分区。
- 拒绝不可用的 ID。
- 按分区身份对当前窗口去重。
- 与已见身份进行比较。
- 应用首次运行和历史回填策略。
- 生成通知候选项。
- 持久化每个有效的已观察 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 调用变成监控器之前,我都会回答这些问题:
- 稳定身份是什么? 必要时按来源、租户、渠道或账号进行分区。
- 哪些字段是可变的? 除非变更检测本身就是目标,否则将它们排除在身份之外。
- 基线存放在哪里? 让某一个层级拥有明确的所有权。
- 首次非空运行时会发生什么? 建立基线或发出告警——绝不能是意外行为。
- 状态缺失或损坏时会发生什么? 故障关闭,而不是重放历史。
- 轮询窗口会饱和吗? 在去重之前先暴露这一状况。
- 缺失意味着什么? 在受限窗口中,缺失并不证明删除。
- 执行会重叠吗? 如果会,使用具有明确并发保证的存储。
一个定时抓取器只会反复告诉你当前存在什么。一个监控器则告诉你什么发生了变化,同时不会编造确定性、重放旧记录或掩盖缺口。
差异不在于多一次 API 调用,而在于状态契约。
资源
- Telegram Public Channel Monitor Actor
- Apify MCP server documentation
- Apify synchronous Actor endpoint
- n8n workflow static data documentation
- Telegram Bot API: Message
- Tested companion repository
_AI 辅助披露:我在起草、编辑和 QA 过程中使用了 AI 工具。我已审阅最终文章及支撑证据,并对其主张和结论负责。_
原文:https://dev.to/apify/a-scraper-isnt-a-monitor-detecting-new-telegram-posts-with-apify-and-n8n-31f5(作者 @bhagyarathnasekara)



