基于Telethon与Flask构建Telegram关键词监听机器人:实现自动化信息监控与实时响应
2026/9/21 11:26:25 网站建设 项目流程

简介:本资源是一套面向PHP开发者与安全研究者的Telegram群组关键词监听机器人源码,聚焦个人学习场景,帮助理解即时通讯平台中隐蔽式内容监控的技术实现原理。系统采用普号(普通账号)伪装策略,通过多账户潜伏于群组,实时捕获含预设关键词的消息,并触发人工响应流程,具备高隐匿性与低感知度特点,适用于安全攻防演练、协议分析及自动化消息处理机制研究。压缩包共57个文件,以36个PHP核心逻辑文件(含事件监听、日志记录、API封装、Swoole/Amp迁移适配等模块)为主,辅以9个备份文件(.zbak)、2个配置JSON、1个Docker编排文件及启动/重载/停止脚本等,整体仅94KB,轻量易部署。目前已有125人学习下载,提供完整可运行架构、清晰分层的Observer模式设计、环境适配脚本及中文教程文档,便于快速理解消息监听链路、事件分发机制与Telegram Bot API集成要点。

1. 项目缘起:从“信息焦虑”到“主动监听”的转变

做社群运营或者项目管理的朋友,大概都经历过这种场景:你负责的某个重要社群,每天消息成千上万,你不可能24小时盯着屏幕。但你又生怕错过某个关键信息,比如客户的重要反馈、竞品的最新动态、或者某个突发事件的讨论。这种“信息焦虑”在快节奏的线上协作中尤为突出。传统的做法是设置消息免打扰,然后定时爬楼,效率低下且容易遗漏。更高级一点的做法是依赖群管理员的“人肉”提醒,但这又增加了人力成本,且无法保证及时性。

正是在这种需求背景下,“关键词监听机器人”的概念应运而生。它本质上是一个自动化信息过滤器,能够7x24小时驻留在指定的社群中,像一位不知疲倦的哨兵,安静地扫描每一条新消息。一旦发现预设的关键词组合,它就能立即触发一系列动作:通知你、记录到数据库、甚至自动回复或转发到其他协作平台。这不仅仅是简单的“@全体成员”,而是一种精准、静默、高效的“信息雷达”。

我这次分享的,就是基于Telegram平台(以下简称TG)构建这样一套系统的完整思路与核心源码实现。它不仅仅是一个简单的“机器人”,而是一个包含“普号隐身监控”与“实时人工响应”的轻量级系统。所谓“普号”,指的是使用普通的个人Telegram账号(而非Bot API创建的机器人账号)作为监听终端,这带来了更高的隐蔽性和灵活性。“隐身监控”意味着机器人读取消息但尽量不发言、不改变群组状态,降低被察觉的风险。而“实时人工响应”则是在机器人捕获到关键信息后,通过接口即时通知到真人,由真人进行决策和回复,实现了“机筛人判”的协同模式。

这套方案特别适合需要从公开或半公开的Telegram群组中追踪特定话题、竞品动态、用户反馈,或者进行舆情监控的团队。接下来,我将从技术选型、核心原理、代码实现、部署细节以及我踩过的几个大坑,为你完整拆解这个项目。

2. 技术栈选型与核心原理:为什么是Telethon+Flask

构建一个TG机器人,官方首推的是通过@BotFather创建并使用Bot API。这种方式稳定、合规,但有明显限制:Bot无法读取普通群组的消息历史,在需要“隐身”监听的大多数群组中基本无用武之地。因此,要实现真正的“普号监控”,我们必须使用用户账号(User Account)来模拟客户端登录。这就是Telethon库的核心价值所在。

2.1 为什么选择Telethon?

Telethon是一个强大且底层的Python MTProto库。MTProto是Telegram客户端与服务器通信的私有协议。使用Telethon,我们可以用编程方式控制一个真实的Telegram用户账号,执行几乎所有手机或桌面客户端能做的操作:登录、收发消息、获取群列表、读取历史消息等等。相较于更高层封装的python-telegram-bot(仅支持Bot API),Telethon给了我们潜入“深海”的能力。

它的另一个巨大优势是异步(asyncio)原生支持。Telegram的通信是高度异步的,使用异步框架可以让我们用单线程高效地处理多个群组的消息流,而不必担心阻塞。这对于需要同时监控数十个群的场景至关重要。

2.2 后端框架:轻量级的Flask

监听机器人本身是“事件驱动”的,它持续运行,等待消息事件。但我们需要一个控制面板和消息推送接口。这就是引入Flask的原因。Flask作为一个轻量级Web框架,可以快速搭建几个API端点:

  • 管理端点:用于动态添加/删除监控的关键词、群组,查看监控日志。
  • 回调端点(Webhook):当Telethon客户端监听到关键词时,将消息详情通过HTTP POST请求发送到这个端点,从而触发后续的“实时人工响应”流程,比如发送通知到钉钉、飞书或企业内部系统。

这种Telethon(异步事件监听) +Flask(同步HTTP服务)的架构,是一种经典的生产者-消费者模式。两者通常运行在同一个进程的不同线程或通过消息队列(如Redis)解耦。在本项目的初始版本中,为了简化部署,我采用了线程内通信的方式。

2.3 核心工作流程

整个系统的运行流程可以概括为以下几步:

  1. 初始化与登录:使用Telethon创建一个TelegramClient实例,填入从Telegram官网申请到的api_idapi_hash。程序首次运行会要求输入手机号验证码,登录成功后会话信息会保存在本地.session文件,后续启动无需重复验证。
  2. 加载监听任务:从数据库或配置文件中读取需要监控的群组ID(或用户名)和对应的关键词列表。
  3. 注册事件处理器:为TelegramClient注册on.NewMessage事件监听器。每当其登录账号所在群组有新消息时,这个回调函数就会被触发。
  4. 消息过滤与处理:在事件处理器中,检查新消息的chat_id是否在监控列表内,并检查消息文本是否包含任何预设关键词(支持简单的模糊匹配或正则表达式)。如果匹配,则执行核心动作。
  5. 触发响应:核心动作包括:将消息详情(发送者、时间、内容、原始消息ID等)格式化,然后通过requests库同步调用本地Flask服务器提供的Webhook接口。
  6. 人工响应闭环:Flask的Webhook接口收到数据后,可以将其推送至办公IM(如通过钉钉机器人、飞书Webhook),提醒相关责任人。责任人点击通知链接,可以快速跳转到TG群对应的消息位置,进行人工回复。

注意:使用用户账号进行自动化操作,必须严格遵守Telegram的服务条款。过度频繁的消息获取或发送行为可能导致账号被暂时限制。因此,代码中必须加入适当的延迟(例如asyncio.sleep)和错误处理,模拟人类操作节奏,这是实现“隐身”的关键技术点之一。

3. 核心代码拆解:从登录到消息推送

下面,我将分模块展示最核心的代码片段,并解释关键逻辑。假设我们的项目结构如下:

tg_keyword_monitor/ ├── config.py # 配置文件 ├── monitor.py # Telethon监听主程序 ├── web_server.py # Flask Web服务器 ├── database.py # 数据库模型(可选,这里用文件代替) └── requirements.txt # 依赖包

3.1 配置文件 (config.py)这里存放敏感信息和全局配置。切记不要将api_id,api_hash.session文件提交到公开仓库!

# config.py import os # 从环境变量读取,更安全 API_ID = int(os.getenv('TG_API_ID', '你的_api_id')) API_HASH = os.getenv('TG_API_HASH', '你的_api_hash') PHONE_NUMBER = os.getenv('TG_PHONE', '+861234567890') # 绑定的手机号 # Flask服务器配置 WEBHOOK_HOST = '127.0.0.1' WEBHOOK_PORT = 5000 WEBHOOK_URL = f'http://{WEBHOOK_HOST}:{WEBHOOK_PORT}/webhook' # 监听配置(实际应从数据库读取) MONITOR_CONFIG = { # 群组ID或用户名: [关键词列表] -1001234567890: ['bug', '故障', '无法登录', 'error'], # 一个技术群 'group_username': ['价格', '优惠', '打折', '竞品名称'], # 一个公开群 } # 请求间隔,避免风控 REQUEST_DELAY = 1.0 # 秒

3.2 Telethon监听核心 (monitor.py)这是系统的大脑,负责登录、监听和初步过滤。

# monitor.py import asyncio import logging from telethon import TelegramClient, events from telethon.tl.types import PeerChannel, PeerChat import requests import json from config import API_ID, API_HASH, PHONE_NUMBER, WEBHOOK_URL, MONITOR_CONFIG, REQUEST_DELAY logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) client = TelegramClient(PHONE_NUMBER, API_ID, API_HASH) def send_to_webhook(chat_title, sender_name, message_text, message_link): """将匹配到的消息发送到Flask Webhook""" payload = { 'chat_title': chat_title, 'sender': sender_name, 'message': message_text, 'link': message_link, 'timestamp': asyncio.get_event_loop().time() } try: resp = requests.post(WEBHOOK_URL, json=payload, timeout=5) if resp.status_code == 200: logger.info(f"Webhook推送成功: {chat_title} - {sender_name}") else: logger.error(f"Webhook推送失败: {resp.status_code}") except Exception as e: logger.error(f"发送Webhook请求异常: {e}") @client.on(events.NewMessage()) async def keyword_monitor(event): """核心事件处理器:过滤并处理新消息""" # 1. 获取聊天信息 chat = await event.get_chat() chat_id = event.chat_id # 2. 检查是否为目标监控群组 if chat_id not in MONITOR_CONFIG: return # 3. 获取消息文本(处理纯文本和带链接/格式的消息) message_text = event.raw_text or "" if not message_text: # 可能是图片、文件等,可以根据需要处理,这里只处理文本 return # 4. 关键词匹配 keywords = MONITOR_CONFIG[chat_id] matched_keywords = [] for kw in keywords: # 简单的大小写不敏感包含匹配,可升级为正则 if kw.lower() in message_text.lower(): matched_keywords.append(kw) if not matched_keywords: return # 5. 匹配成功,准备数据 sender = await event.get_sender() sender_name = f"{sender.first_name or ''} {sender.last_name or ''}".strip() or sender.username or '未知用户' chat_title = chat.title if hasattr(chat, 'title') else '私聊/频道' # 构造消息链接(方便人工快速跳转) if hasattr(chat, 'username') and chat.username: message_link = f"https://t.me/{chat.username}/{event.id}" else: # 对于私有群组,链接无法直接跳转,记录ID message_link = f"chat_id: {chat_id}, msg_id: {event.id}" logger.info(f"[命中] 群组「{chat_title}」- 用户「{sender_name}」- 关键词「{matched_keywords}」") logger.info(f" 内容: {message_text[:100]}...") # 6. 触发Webhook,推送消息详情 # 使用asyncio.to_thread将同步的requests调用放到线程池执行,避免阻塞异步循环 await asyncio.to_thread(send_to_webhook, chat_title, sender_name, message_text, message_link) # 7. 礼貌性延迟,模拟人类,降低请求频率 await asyncio.sleep(REQUEST_DELAY) async def main(): """主异步函数""" await client.start(phone=PHONE_NUMBER) logger.info("监听机器人启动成功!") # 打印已加入的群组,方便配置 async for dialog in client.iter_dialogs(): if dialog.is_group or dialog.is_channel: logger.info(f"群组/频道: {dialog.name} (ID: {dialog.id})") # 持续运行,直到接收到停止信号 await client.run_until_disconnected() if __name__ == '__main__': # 启动异步事件循环 with client: client.loop.run_until_complete(main())

3.3 Flask Web服务器 (web_server.py)这是一个简单的HTTP服务器,接收监听器的推送并转发。

# web_server.py from flask import Flask, request, jsonify import logging import json app = Flask(__name__) logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) # 这里可以集成钉钉、飞书、企业微信等机器人的推送函数 def send_to_dingtalk(webhook_url, content): """示例:推送消息到钉钉群""" import requests headers = {'Content-Type': 'application/json'} data = { "msgtype": "markdown", "markdown": { "title": "TG关键词监控告警", "text": content } } try: resp = requests.post(webhook_url, headers=headers, data=json.dumps(data)) return resp.status_code == 200 except Exception as e: logger.error(f"钉钉推送失败: {e}") return False @app.route('/webhook', methods=['POST']) def handle_webhook(): """接收Telethon监听器推送的Webhook""" if not request.is_json: return jsonify({'error': 'Invalid content type'}), 400 data = request.get_json() chat_title = data.get('chat_title', 'N/A') sender = data.get('sender', 'N/A') message = data.get('message', 'N/A') link = data.get('link', '#') logger.info(f"收到告警 -> 来源: {chat_title}, 发送者: {sender}") # 构建推送内容 markdown_content = f"""### TG监控告警 **群组**: {chat_title} **发送者**: {sender} **关键词消息**: > {message[:200]} **快速跳转**: [查看原文]({link}) """ # 调用推送函数(这里以钉钉为例,需配置自己的Webhook) dingtalk_webhook = "https://oapi.dingtalk.com/robot/send?access_token=YOUR_TOKEN" success = send_to_dingtalk(dingtalk_webhook, markdown_content) if success: return jsonify({'status': 'ok', 'msg': '推送成功'}), 200 else: return jsonify({'status': 'error', 'msg': '推送失败'}), 500 @app.route('/') def index(): return "TG关键词监控系统 Webhook 服务运行中。" if __name__ == '__main__': # 注意:在生产环境中,不要使用debug模式,并用WSGI服务器(如Gunicorn)运行 app.run(host='0.0.0.0', port=5000, debug=False)

4. 部署实战与隐身技巧:让机器人“活”在背景里

代码写好了,但让这个机器人稳定、隐蔽地长期运行,才是真正的挑战。直接在本机运行python monitor.pypython web_server.py不是长久之计。

4.1 进程管理与持久化

推荐使用systemd(Linux)或Supervisor来管理这两个进程。以systemd为例,创建两个服务文件:

/etc/systemd/system/tg-monitor.service

[Unit] Description=TG Keyword Monitor Daemon After=network.target [Service] Type=simple User=your_username WorkingDirectory=/path/to/tg_keyword_monitor Environment="PATH=/usr/local/bin:/usr/bin" ExecStart=/usr/bin/python3 /path/to/tg_keyword_monitor/monitor.py Restart=always RestartSec=10 StandardOutput=syslog StandardError=syslog SyslogIdentifier=tg-monitor [Install] WantedBy=multi-user.target

/etc/systemd/system/tg-webhook.service

[Unit] Description=TG Monitor Webhook Server After=network.target [Service] Type=simple User=your_username WorkingDirectory=/path/to/tg_keyword_monitor Environment="PATH=/usr/local/bin:/usr/bin" ExecStart=/usr/bin/python3 /path/to/tg_keyword_monitor/web_server.py Restart=always RestartSec=10 StandardOutput=syslog StandardError=syslog SyslogIdentifier=tg-webhook [Install] WantedBy=multi-user.target

然后使用sudo systemctl daemon-reloadsudo systemctl start tg-monitor tg-webhook启动,并sudo systemctl enable设置开机自启。这样即使服务器重启,服务也会自动恢复。

4.2 关键的“隐身”配置与风控规避

使用普通账号进行自动化操作,最大的风险就是被Telegram风控系统识别并限制。以下是我在实践中总结出的几条“军规”:

  1. 会话管理Telethon生成的.session文件是登录凭证。务必妥善保管,并设置正确的文件权限(如600),防止泄露。一个.session文件理论上可以永久使用,除非你在其他地方登录此账号导致其失效。
  2. 速率限制:代码中的REQUEST_DELAY至关重要。不要设置得太小(如低于0.5秒)。对于消息事件处理函数,即使匹配到关键词,也最好在函数末尾加一个短暂的sleepTelethon内部有自动的洪水等待(Flood Wait)处理,但我们主动放缓节奏是更友好的行为。
  3. 避免敏感操作:监听机器人绝对不要在监控的群组内主动发言、点赞、转发或修改群设置。它的行为模式应尽可能接近一个“只读”的隐身用户。如果需要测试,请在自己的私人小群进行。
  4. IP地址稳定性:尽量在固定的、干净的IP地址(如云服务器IP)下运行机器人。频繁切换IP(尤其是数据中心IP和家庭IP混用)容易触发安全警报。
  5. 模拟人类行为:可以随机化延迟时间,比如await asyncio.sleep(1 + random.random())。在启动时,可以模拟人类偶尔翻看历史消息的行为(使用client.get_messages),但频率要低。
  6. 准备备用号:重要项目不要只依赖一个监控账号。可以准备2-3个备用号,使用类似的代码但不同的.session文件运行,监控相同的群组。即使一个号被限制,其他号也能顶上。

4.3 监控配置的动态化管理

上面的示例将配置硬编码在config.py中,这不利于维护。生产环境应该使用数据库(如SQLite或PostgreSQL)来管理监控任务。

可以设计一张表monitor_tasks

CREATE TABLE monitor_tasks ( id INTEGER PRIMARY KEY AUTOINCREMENT, chat_id BIGINT NOT NULL, -- 群组ID chat_name TEXT, -- 群组名称(便于管理) keywords TEXT NOT NULL, -- JSON数组,如 '["bug", "error"]' is_active BOOLEAN DEFAULT 1, -- 是否启用 created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );

然后在monitor.py中,定期(例如每5分钟)从数据库拉取最新的有效任务列表,更新内存中的MONITOR_CONFIG。Flask服务也可以提供简单的管理页面,来增删改查这些监控任务。

5. 踩坑实录与进阶优化

在开发和运行这个系统的过程中,我遇到了不少问题,这里分享几个典型的坑和解决方案。

5.1 坑一:api_idapi_hash申请与使用

很多人第一步就卡在这里。api_idapi_hash不是Bot Token,它们代表一个“应用”,用于用户账号登录。

  • 步骤:访问 https://my.telegram.org ,用你的手机号登录。进入“API development tools”,填写应用信息(随便填,如MyMonitor),提交后即可获得api_idapi_hash
  • 坑点:同一个IP频繁申请api_id可能会被限制。api_hash务必保密,泄露可能导致账号安全风险。

5.2 坑二:无法获取私有群组的chat_id或监听不到消息

这是最常见的问题。

  • 原因与解决:你的监听账号必须已经是该群组的成员。对于公开群,你可以用client.get_entity(‘群组用户名’)来解析。对于私有群,最稳妥的方式是先用账号人工加入群组,然后在monitor.pymain函数启动时,它会遍历并打印所有对话的ID,把这个ID记录下来填入配置。
  • 权限问题:即使你在群里,如果是被禁言状态或群组设置了严格的权限,也可能无法读取消息。确保账号在群内有基本的“查看消息”权限。

5.3 坑三:asyncio事件循环冲突

如果你在Flask这样的同步Web框架中直接调用Telethon的异步函数,或者尝试在Jupyter Notebook中运行,经常会遇到Event loop is closedEvent loop already running的错误。

  • 根本原因asyncio不允许在同一个线程中嵌套运行事件循环。
  • 解决方案
    1. 分离进程:就像本项目设计的一样,将异步的监听程序(monitor.py)和同步的Web服务器(web_server.py)作为两个独立的进程运行。这是最清晰、最稳定的方式。
    2. 使用asyncio.run():确保你的异步主函数被asyncio.run(main())调用。在复杂项目中,可以使用nest_asyncio库修补,但这通常是下策,可能引入不确定性。
    3. 线程隔离:如果必须在同步代码中偶尔调用异步客户端(例如在Flask路由里手动抓取一次消息),可以使用asyncio.run_coroutine_threadsafe并配合一个全局的、长期运行的事件循环线程。

5.4 进阶优化方向

  1. 关键词策略升级:目前的简单包含匹配误报率高。可以引入:
    • 正则表达式:实现更复杂的模式匹配,如匹配特定格式的版本号v\d+\.\d+
    • 分词与语义过滤:结合jieba(中文)等分词库,避免子串误匹配(如“苹果”匹配到“苹果手机”是合理的,但“代码”匹配到“密码”就不合理)。
    • 负面词过滤:定义“排除词”列表,当消息同时包含关键词和排除词时忽略。
  2. 消息上下文获取:有时单条消息意义不明,需要上下文。可以在匹配到关键词后,用client.get_messages(chat_id, min_id=msg_id-5, max_id=msg_id+5)来获取附近的消息,一并推送给人工判断。
  3. 消息去重:同一个问题可能在群内被多人反复提及。可以基于消息内容的哈希值或相似度,在一段时间内(如10分钟)进行去重,避免轰炸通知。
  4. 状态持久化与高可用:将匹配记录、发送状态存入数据库。即使服务重启,也能知道哪些消息已处理。可以考虑使用Redis作为消息队列,将Telethon监听器作为生产者,将推送逻辑作为多个消费者,实现解耦和水平扩展。
  5. 容器化部署:使用Docker封装整个应用,将配置、.session文件通过卷挂载,可以极大简化部署和迁移流程。

构建这样一个系统,技术本身并不复杂,真正的价值在于如何将它无缝融入你的工作流,并持续稳定地运行。它解放了你的双眼,让你从海量的信息噪音中抽身,专注于那些真正需要你回应的信号。从“人找信息”到“信息找人”,这小小的改变,带来的效率提升是巨大的。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询