
简介本资源是一套面向PHP开发者与安全研究者的Telegram群组关键词监听机器人源码聚焦个人学习场景帮助理解即时通讯平台中隐蔽式内容监控的技术实现原理。系统采用普号普通账号伪装策略通过多账户潜伏于群组实时捕获含预设关键词的消息并触发人工响应流程具备高隐匿性与低感知度特点适用于安全攻防演练、协议分析及自动化消息处理机制研究。压缩包共57个文件以36个PHP核心逻辑文件含事件监听、日志记录、API封装、Swoole/Amp迁移适配等模块为主辅以9个备份文件.zbak、2个配置JSON、1个Docker编排文件及启动/重载/停止脚本等整体仅94KB轻量易部署。目前已有125人学习下载提供完整可运行架构、清晰分层的Observer模式设计、环境适配脚本及中文教程文档便于快速理解消息监听链路、事件分发机制与Telegram Bot API集成要点。1. 项目缘起从被动等待到主动感知的业务需求做社群运营或者信息监控的朋友应该都遇到过类似的痛点你加入了一个几百甚至上千人的大群里面信息刷得飞快你不可能24小时盯着屏幕。但你又生怕错过那些提到你公司、产品、竞品或者特定关键词的重要消息。一旦错过可能就是一次商机流失或者一次负面舆情发酵的起点。传统的做法要么是手动爬楼效率低下要么是依赖一些功能简单的机器人只能做到关键词匹配后你一下你依然需要时刻关注手机通知本质上还是“人找信息”。我手头这个项目就是为了解决这个核心矛盾而诞生的。它不是一个简单的“关键词提醒”工具而是一套完整的“普号隐身监控与实时人工响应系统”。这个名字听起来有点拗口拆解一下“普号”指的是普通的、无特殊权限的账号“隐身监控”意味着这个监控行为对群内其他成员是近乎不可见的“实时人工响应”则是在机器捕获到关键信息后能无缝衔接人工介入进行回复、互动或记录。这套系统的价值在于它将人力从枯燥的“盯屏”中解放出来转而投入到更有价值的“决策与交互”环节实现了从“人盯信息”到“信息找人人处理信息”的范式转变。最近在相关社群里关于“TG机器人”、“监控”、“源码”的讨论热度一直很高大家的需求很明确但苦于找不到一套稳定、隐蔽且功能完整的解决方案。很多现成的机器人要么功能单一要么容易被风控要么就是源码晦涩难以二次开发。因此我决定将之前为某个海外市场运营团队搭建的这套系统进行脱敏和重构把核心思路和关键代码分享出来。本文将深入拆解其设计原理、实现细节以及那些在实战中积累下来的、关乎稳定性和隐蔽性的宝贵经验。2. 系统架构设计兼顾隐匿、稳定与可扩展在开始敲代码之前我们必须先想清楚整个系统应该如何运作。一个鲁棒的监控响应系统绝不能是简单写个脚本就了事。它需要分层设计各司其职以应对网络波动、平台风控、业务变更等各种挑战。2.1 核心组件与数据流整个系统可以划分为四个核心层数据像流水线一样在其中传递数据采集层Client这是系统的“眼睛”和“耳朵”。由一个或多个普通账号即“普号”担任。它们通过官方API例如Telegram Bot API的getUpdates轮询或更推荐的MTProto协议库安静地潜伏在目标群组中接收所有流经的消息。这一层的关键词是“低调”和“稳定”要模拟真人客户端的正常行为避免频繁请求或异常动作触发平台的风控机制。消息过滤与路由层Filter Router这是系统的“大脑皮层”。采集层获取的原始消息是海量且杂乱的。这一层负责进行初步清洗和关键判断。首先它会过滤掉系统消息、服务消息或无意义的更新。然后将剩余的消息文本与预先配置的关键词规则库进行匹配。这里的规则不仅仅是简单的字符串包含还应支持正则表达式、模糊匹配甚至简单的意图识别例如匹配“价格”且上下文情绪为“负面”。一旦匹配成功该条消息连同其上下文前几条消息、发送者信息、时间戳等元数据会被打包成一个“事件”并路由到下一层。人工响应接口层Interface这是连接机器与人的“桥梁”。当过滤层产生一个高优先级的事件时系统需要以某种方式通知真人操作员。最直接的方式是通过另一个Telegram Bot向指定的管理员账号发送警报。但更优雅的做法是集成到一个内部工作台Web Dashboard上操作员可以在网页上直接看到警报列表、消息原文、上下文并无需切换应用就能进行回复、打标签或标记为已处理。这极大地提升了响应效率。持久化与审计层Storage Audit这是系统的“记忆”。所有原始消息、匹配到的事件、操作员的所有响应动作都需要被安全地存储下来。这不仅是出于审计和复盘的需求例如分析舆情趋势更是系统稳定性的保障。万一监控服务短暂中断持久化队列可以确保消息不丢失恢复后能继续处理。通常我们会使用如 PostgreSQL 或 MySQL 存储结构化数据事件、用户用 Redis 作为高速缓存和消息队列处理实时推送而原始消息文本可能还会存一份到 Elasticsearch 以便后续全文检索。[ 目标群组 ] - [ 普号客户端 (MTProto) ] - [ 消息过滤/路由服务 ] - [ 消息队列 (Redis) ] | v [ 管理员 ] - [ Web管理后台 / Bot警报 ] - [ 人工响应服务 ] - [ 数据库 (PostgreSQL) ]2.2 为什么选择“普号”而非“Bot”作为采集端这是一个关键设计决策。Telegram 官方提供了功能强大的 Bot API用它来接收群组消息非常简单前提是Bot被加入群组且赋予了相应权限。但Bot账号有几个致命缺点显眼性Bot在群成员列表里有特殊标识它的存在本身就是公开的。功能限制Bot无法接收群组中所有成员的普通消息除非消息以“/”命令开头或直接回复了Bot。虽然可以通过设置privacy mode为False来接收所有消息但这并非所有群组都允许且可能引起群主或成员反感。风控敏感官方对Bot的行为监控更严格频繁请求或异常模式容易被限流甚至封禁。而使用普通用户账号通过模拟MTProto协议则完全模拟了一个真实用户的在线行为高度隐蔽在成员列表中就是一个普通用户没有任何特殊标记。信息完整可以接收到群内所有的消息、媒体文件、甚至已删除消息的提示取决于客户端实现。行为自然可以控制上线时间、读取状态已读回执需谨慎使用更不容易被平台的风控系统识别为机器人。当然这带来了更高的技术复杂度需要处理MTProto协议、会话授权登录、2FA验证等。但为了长期的稳定和隐蔽这个代价是值得的。市面上成熟的库如Telethon(Python) 或MadelineProto(PHP) 已经很好地封装了这些复杂性。3. 关键技术实现从登录监听到的完整链路接下来我们深入到代码层面看看各个核心环节如何实现。这里以 Python 生态为例使用Telethon库作为MTProto客户端。3.1 普号客户端的初始化与安全登录第一步是让我们的“普号”安全上线。绝对不要将账号、密码或API密钥硬编码在源码中。# config.py - 配置文件 import os from dotenv import load_dotenv load_dotenv() # 从 .env 文件加载环境变量 API_ID int(os.getenv(TG_API_ID)) # 从 my.telegram.org 申请 API_HASH os.getenv(TG_API_HASH) PHONE_NUMBER os.getenv(TG_PHONE) # 国际格式如 8613012345678 SESSION_NAME monitor_session # 会话文件名称 # 目标群组的用户名或ID列表 TARGET_GROUPS [group_username_1, -1001234567890] # 关键词列表支持正则表达式 KEYWORDS [ r(?i)竞品名称A|竞品名称B, # 不区分大小写匹配竞品 r价格\s*(过高|太贵|上涨), # 匹配价格负面词汇 rbug|故障|无法使用, # 匹配问题反馈 r合作|询盘|contact us, # 匹配商机 ]# client.py - 客户端初始化与登录 from telethon import TelegramClient, events from telethon.tl.types import PeerChannel, PeerChat import asyncio import logging from config import API_ID, API_HASH, PHONE_NUMBER, SESSION_NAME logging.basicConfig(levellogging.INFO) logger logging.getLogger(__name__) class MonitorClient: def __init__(self): self.client TelegramClient(SESSION_NAME, API_ID, API_HASH) # 连接状态标志 self.connected False async def connect(self): 连接并登录客户端 try: await self.client.connect() if not await self.client.is_user_authorized(): logger.info(f首次登录账号 {PHONE_NUMBER}...) await self.client.send_code_request(PHONE_NUMBER) # 这里需要手动输入收到的验证码 code input(请输入Telegram发送的验证码: ) await self.client.sign_in(PHONE_NUMBER, code) # 如果设置了2FA密码还需要以下步骤 # password input(请输入2FA密码: ) # await self.client.sign_in(passwordpassword) else: logger.info(会话有效已自动登录。) self.connected True logger.info(客户端登录成功) except Exception as e: logger.error(f连接或登录失败: {e}) self.connected False raise async def start_monitoring(self, handler_callback): 开始监听消息并注册事件处理器 if not self.connected: await self.connect() self.client.on(events.NewMessage()) async def message_handler(event): # 将事件传递给外部定义的处理回调函数 await handler_callback(event) logger.info(开始监听消息...) await self.client.run_until_disconnected()关键经验1会话管理与2FA处理Telethon会将登录会话保存在本地文件如monitor_session.session中。首次登录后后续启动通常无需再输入验证码除非会话失效。务必妥善保管这个.session文件它等同于账号的登录凭证。如果账号启用了两步验证2FA代码中需要有相应的分支来处理密码输入。在生产环境中可以通过环境变量或安全的密钥管理服务来传递2FA密码但更推荐使用专门为自动化创建的、未启用2FA的账号。3.2 智能消息过滤与事件生成消息处理器是核心逻辑所在。我们需要判断一条消息是否值得关注并为其添加上下文。# filter_engine.py - 过滤引擎 import re from datetime import datetime import asyncio from typing import List, Optional, Dict, Any class KeywordFilterEngine: def __init__(self, keyword_patterns: List[str]): 初始化过滤引擎 :param keyword_patterns: 关键词正则表达式列表 self.compiled_patterns [re.compile(pattern, re.IGNORECASE) for pattern in keyword_patterns] self.message_cache {} # 用于缓存上下文消息键为对话ID async def process_message(self, event) - Optional[Dict[str, Any]]: 处理单条消息返回匹配到的事件字典否则返回None msg event.message # 1. 基础过滤忽略空消息、服务消息、自己发送的消息 if not msg.message or msg.out: return None # 获取对话信息 chat await event.get_chat() chat_id event.chat_id chat_title getattr(chat, title, Private Chat) # 2. 判断是否为目标群组可从配置读取这里简化演示 # if chat_id not in TARGET_GROUPS: return None # 3. 关键词匹配 matched_keywords [] for pattern in self.compiled_patterns: if pattern.search(msg.message): matched_keywords.append(pattern.pattern) if not matched_keywords: return None # 没有匹配到任何关键词 # 4. 获取消息上下文例如前2条消息 context_messages await self._get_context(chat_id, msg.id, limit2) # 5. 构建事件对象 event_data { event_id: f{chat_id}_{msg.id}_{int(datetime.now().timestamp())}, chat_id: chat_id, chat_title: chat_title, message_id: msg.id, sender_id: msg.sender_id, sender_name: self._get_sender_name(msg), # 需要实现获取发送者名字的函数 text: msg.message, timestamp: msg.date.isoformat(), matched_keywords: matched_keywords, context: context_messages, raw_event: event.to_dict() # 保存原始事件数据以备不时之需 } logger.info(f检测到关键词事件于 [{chat_title}]: {matched_keywords}) return event_data async def _get_context(self, chat_id, message_id, limit2): 获取某条消息之前的几条消息作为上下文 # 这里需要根据使用的客户端库实现获取历史消息的逻辑 # 例如使用 client.get_messages context [] # 伪代码实际实现需根据Telethon API调整 # try: # async for msg in client.iter_messages(chat_id, limitlimit, offset_idmessage_id, reverseTrue): # if msg.id message_id: # context.append({id: msg.id, text: msg.text, sender: msg.sender_id}) # except Exception as e: # logger.error(f获取上下文失败: {e}) return context def _get_sender_name(self, msg): 提取发送者显示名称 sender msg.sender if sender: if getattr(sender, first_name, None) or getattr(sender, last_name, None): return f{getattr(sender, first_name, )} {getattr(sender, last_name, )}.strip() elif getattr(sender, title, None): # 对于频道/群组 return sender.title elif getattr(sender, username, None): return f{sender.username} return Unknown关键经验2匹配策略与性能正则表达式虽然强大但在海量消息实时匹配时需注意性能。过于复杂的正则可能导致CPU尖峰。建议将最常出现、最核心的关键词放在列表前面。对于简单的固定字符串使用in操作符判断可能比正则更快可以先做一层快速过滤。考虑引入AC自动机Aho-Corasick算法来处理大量关键词的匹配效率远高于遍历正则列表。Python有ahocorasick库可以实现。3.3 事件队列与人工响应桥接过滤后的事件不能直接处理应该先放入一个消息队列实现生产与消费的解耦提高系统的抗压能力和可扩展性。# queue_manager.py - 基于Redis的简单事件队列 import json import redis import logging from typing import Dict, Any logger logging.getLogger(__name__) class EventQueueManager: def __init__(self, redis_hostlocalhost, redis_port6379, queue_nametg_monitor_events): self.redis_client redis.Redis(hostredis_host, portredis_port, decode_responsesTrue) self.queue_name queue_name try: self.redis_client.ping() logger.info(Redis连接成功。) except redis.ConnectionError: logger.error(无法连接到Redis请检查服务是否运行。) raise async def push_event(self, event_data: Dict[str, Any]): 将事件推入Redis队列 try: event_json json.dumps(event_data, ensure_asciiFalse) # 使用LPUSH将事件放入列表左侧 result self.redis_client.lpush(self.queue_name, event_json) logger.debug(f事件已入队队列长度: {result}) except Exception as e: logger.error(f事件入队失败: {e}) async def pop_event(self, timeout0): 从Redis队列阻塞弹出事件BRPOP try: # BRPOP 是阻塞式弹出timeout0表示无限等待 _, event_json self.redis_client.brpop(self.queue_name, timeouttimeout) if event_json: return json.loads(event_json) except Exception as e: logger.error(f从队列弹出事件失败: {e}) return None有了队列我们就可以创建独立的后端服务消费者来处理这些事件比如发送警报到管理员的Telegram Bot。# alert_bot.py - 警报Bot消费者服务 from telegram import Bot, Update from telegram.ext import Application, CommandHandler, MessageHandler, filters, CallbackContext import asyncio from queue_manager import EventQueueManager import logging logging.basicConfig(format%(asctime)s - %(name)s - %(levelname)s - %(message)s, levellogging.INFO) logger logging.getLogger(__name__) class AlertBot: def __init__(self, token: str, admin_chat_id: int, event_queue: EventQueueManager): self.application Application.builder().token(token).build() self.admin_chat_id admin_chat_id self.event_queue event_queue # 注册命令处理器 self.application.add_handler(CommandHandler(start, self.start)) # 可以添加更多命令... async def start(self, update: Update, context: CallbackContext): await update.message.reply_text(监控警报Bot已启动。) async def consume_and_alert(self): 消费事件队列并发送警报 logger.info(开始消费事件队列...) while True: event await self.event_queue.pop_event(timeout30) # 阻塞30秒 if event: alert_text self._format_alert(event) try: await self.application.bot.send_message( chat_idself.admin_chat_id, textalert_text, parse_modeHTML ) logger.info(f已向管理员发送警报: {event.get(event_id)}) except Exception as e: logger.error(f发送警报失败: {e}) # 短暂休眠避免CPU空转 await asyncio.sleep(0.1) def _format_alert(self, event): 将事件格式化为易读的警报消息 chat_title event.get(chat_title, Unknown Chat) sender event.get(sender_name, Unknown User) keywords , .join(event.get(matched_keywords, [])) text_preview event.get(text, )[:200] # 预览前200字符 msg_link fhttps://t.me/c/{str(event[chat_id]).replace(-100, )}/{event[message_id]} if chat_id in event and message_id in event else # return f b关键词监控警报/b b群组/b: {chat_title} b发送者/b: {sender} b匹配关键词/b: code{keywords}/code b消息预览/b: {text_preview} b上下文/b: [需在管理后台查看] b直达链接/b: a href{msg_link}查看消息/a b事件ID/b: code{event.get(event_id)}/code 请及时处理。3.4 Web管理后台响应与审计一体化对于需要快速响应和团队协作的场景一个简单的Web管理后台远比在Telegram里回复Bot方便。这里以 Flask 为例展示一个极简的响应界面。# app.py - Flask管理后台 from flask import Flask, render_template, request, jsonify import sqlite3 import json from datetime import datetime app Flask(__name__) DATABASE monitor_events.db def init_db(): conn sqlite3.connect(DATABASE) c conn.cursor() c.execute(CREATE TABLE IF NOT EXISTS events (id TEXT PRIMARY KEY, chat_title TEXT, sender_name TEXT, text TEXT, matched_keywords TEXT, context TEXT, timestamp DATETIME, status TEXT DEFAULT pending, -- pending, processing, resolved response TEXT, responded_by TEXT, responded_at DATETIME)) conn.commit() conn.close() app.route(/) def index(): 展示待处理事件列表 conn sqlite3.connect(DATABASE) conn.row_factory sqlite3.Row c conn.cursor() c.execute(SELECT * FROM events WHERE statuspending ORDER BY timestamp DESC LIMIT 50) events c.fetchall() conn.close() # 将events转换为字典列表供模板使用 events_list [dict(event) for event in events] return render_template(dashboard.html, eventsevents_list) app.route(/api/event/event_id/respond, methods[POST]) def respond_to_event(event_id): 处理对某个事件的响应 data request.json response_text data.get(response) responder data.get(responder, admin) # 应从会话中获取真实用户名 if not response_text: return jsonify({error: 回复内容不能为空}), 400 conn sqlite3.connect(DATABASE) c conn.cursor() # 1. 更新事件状态和回复 c.execute(UPDATE events SET statusresolved, response?, responded_by?, responded_at? WHERE id?, (response_text, responder, datetime.now().isoformat(), event_id)) # 2. (可选) 通过Telegram客户端发送回复到原群组 # 这里需要调用之前写的MonitorClient或者通过一个共享的客户端实例发送消息 # await client.send_message(entityevent_chat_id, messageresponse_text, reply_toevent_message_id) conn.commit() conn.close() return jsonify({success: True}) # 一个后台任务用于从Redis队列取出事件并存入数据库 def event_consumer_worker(): # 这里需要连接Redis和数据库循环 pop_event 并插入 pass if __name__ __main__: init_db() # 在生产环境中event_consumer_worker应该在独立进程中运行 app.run(debugTrue)对应的dashboard.html模板可以列出事件并提供快速回复的文本框。这样运营人员在一个网页上就能完成从查看、分析到回复的全流程。4. 部署、风控与实战避坑指南将代码跑起来只是第一步要让这套系统7x24小时稳定、隐蔽地运行还需要在部署和策略上下功夫。4.1 部署架构建议不建议将所有组件客户端、过滤器、队列、Web后台都放在同一台服务器或同一个进程中。分离客户端与后端服务将MonitorClient单独部署在一台或多台VPS上。甚至可以为一个账号部署一个轻量级客户端分散风险。这些客户端只负责监听和推送事件到中央队列Redis。使用进程管理工具使用systemd或supervisor来管理客户端进程实现崩溃后自动重启。为每个客户端编写独立的 service 文件。容器化考虑如果组件多可以考虑使用 Docker Compose 来编排 Redis、数据库、后端消费者和Web服务。但客户端部分由于涉及登录状态.session文件需要持久化存储卷。4.2 应对平台风控的核心策略这是项目能否长期运行的生命线。平台不希望被自动化工具滥用所以我们的行为必须尽可能“像人”。节奏模拟上线延迟启动客户端后不要立即开始疯狂拉取消息。可以随机等待几分钟。消息处理间隔在message_handler中即使收到消息也不一定要立刻处理。可以引入一个随机延迟例如0.5秒到3秒模拟真人阅读时间。主动动作节制除非必要否则监控账号不要在群内发言、点赞或转发。如果为了响应必须发言频率也要极低且内容要自然。连接保活与重连# 在客户端类中增加保活和重连逻辑 async def keep_alive(self): while True: await asyncio.sleep(300) # 每5分钟 try: if not self.client.is_connected(): logger.warning(连接断开尝试重连...) await self.client.connect() # 重连后可能需要重新获取对话列表等状态 except Exception as e: logger.error(f保活/重连检查失败: {e})账号健康度使用“老号”注册时间较长、有过正常聊天记录的账号比新号更安全。避免同一IP地址登录过多账号。定期如每周让账号在手机官方App上登录一下维持活跃度。指纹多样化如果运行多个客户端可以尝试定制不同的device_model、system_version等客户端信息虽然Telethon对此支持有限但尽可能避免完全一致。4.3 常见问题与排查收不到群组消息检查账号是否仍在群内有时管理员会清理不活跃成员。检查会话权限确保.session文件有效且账号未被限制。确认监听代码正确events.NewMessage()会监听所有对话如果加了过滤条件请检查过滤逻辑是否正确。群组类型某些超大群或频道可能需要先调用client.get_participants(chat)或client.get_entity(chat)来“激活”对话列表。出现FloodWaitError 这是触发了平台限流。错误信息中会包含需要等待的秒数。务必在代码中捕获这个异常并休眠指定时间。from telethon.errors import FloodWaitError try: await client.send_message(entity, test) except FloodWaitError as e: logger.warning(f触发限流需要等待 {e.seconds} 秒) await asyncio.sleep(e.seconds 10) # 多等10秒更安全内存或CPU占用过高检查消息处理函数是否阻塞或存在死循环。确认正则匹配是否过于复杂。对于长时间运行的服务确保定期清理内存中的缓存如message_cache。Web后台无法发送回复Flask是同步框架而Telethon客户端是异步的。不能在Flask的同步视图函数中直接await client.send_message。解决方案是使用asyncio.run_coroutine_threadsafe或者将发送消息的任务放入一个异步队列由后台的异步工作线程来消费执行。这套“普号隐身监控与实时人工响应系统”的源码和思路基本涵盖了从底层监听、中间件处理到上层应用的全链路。它不是一个开箱即用的产品而是一个需要根据自身业务需求进行定制和调优的框架。其中关于风控和隐蔽性的经验尤其值得在正式部署前反复思考和测试。技术是为业务服务的在实现自动化监控的同时务必保持对平台规则的尊重和对用户隐私的敬畏将系统用在合规的信息收集与及时的用户服务上才能创造长远的价值。本文还有配套的精品资源点击获取