在开始构建精准流量激活系统之前,我们需要准备好运行环境。本系统基于 Python 3.8+ 开发,利用 Redis 进行高频行为数据的实时存储与计算。请严格按照以下步骤操作,确保环境一致。
1. 安装 Redis 服务
Redis 是本系统的核心,用于处理实时流量计数。如果你已经安装,请跳过此步。
建议下载 Redis Stack 的 Windows 版本 MSI 安装包,直接双击安装并启动服务。下载地址:https://redis.com/download/redis-stack/
2. 创建项目目录与虚拟环境
为了保证依赖隔离,我们创建一个独立的虚拟环境。打开终端执行:
```bash mkdir traffic_activation_system cd traffic_activation_system python3 -m venv venv source venv/bin/activate Windows下使用 venv\Scripts\activate ```3. 安装 Python 依赖库
我们需要 Flask 作为 Web 框架接收流量事件,redis-py 操作数据库。执行以下命令:
```bash pip install flask redis requests ```为了实现“精准”激活,我们不能只做简单的计数。本系统采用滑动窗口算法来判断用户行为。例如:设定规则为“用户在 60 秒内浏览特定页面 3 次”即判定为高意向流量,触发激活。
项目目录结构如下,请手动创建对应文件:
config.py:配置文件,存储 Redis 连接及规则阈值。activation_engine.py:核心引擎,包含滑动窗口算法与触发逻辑。app.py:Web 接口服务,接收前端上报的用户行为。首先编写配置文件,将所有参数集中管理,方便后续调整激活阈值。
```python config.py Redis 配置 REDIS_HOST = 'localhost' REDIS_PORT = 6379 REDIS_DB = 0 REDIS_PASSWORD = None 激活规则配置 规则含义:在 TIME_WINDOW 秒内,达到 THRESHOLD 次行为,则触发激活 ACTIVATION_RULES = { 'view_pricing_page': { 'threshold': 3, 触发次数 'time_window': 60, 时间窗口(秒) 'action': 'send_vip_coupon' 触发后的动作 }, 'click_download_btn': { 'threshold': 1, 'time_window': 10, 'action': 'push_sales_notification' } } ```这是系统的最核心部分。我们使用 Redis 的 Sorted Set (ZSET) 来实现滑动窗口。每一个用户的行为事件都会记录为一个带有时间戳的分数。每次上报行为时,清理掉时间窗口之外的数据,然后统计剩余数量。
```python activation_engine.py import redis import time import config import logging 配置日志 logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') logger = logging.getLogger(__name__) class ActivationEngine: def __init__(self): self.redis_client = redis.StrictRedis( host=config.REDIS_HOST, port=config.REDIS_PORT, db=config.REDIS_DB, password=config.REDIS_PASSWORD, decode_responses=True ) def record_event_and_check(self, user_id, event_type): """ 记录事件并检查是否满足激活条件 :param user_id: 用户唯一标识 :param event_type: 事件类型,需在 config.py 中定义 :return: (bool, str) 第一个元素表示是否触发激活,第二个元素表示触发的动作名称 """ if event_type not in config.ACTIVATION_RULES: logger.warning(f"未定义的事件类型: {event_type}") return False, None rule = config.ACTIVATION_RULES[event_type] threshold = rule['threshold'] time_window = rule['time_window'] action_name = rule['action'] 构造 Redis Key,格式:act:用户ID:事件类型 redis_key = f"act:{user_id}:{event_type}" current_timestamp = time.time() 1. 移除时间窗口之外的数据 (ZREMRANGEBYSCORE) 保留 [当前时间 - 窗口时间, 当前时间] 之间的数据 min_score = current_timestamp - time_window self.redis_client.zremrangebyscore(redis_key, 0, min_score) 2. 添加当前行为记录 (ZADD) member 使用时间戳字符串确保唯一性,score 使用时间戳用于排序 self.redis_client.zadd(redis_key, {str(current_timestamp): current_timestamp}) 3. 设置 Key 过期时间,防止冷数据长期占用内存 (过期时间设为窗口时间的2倍) self.redis_client.expire(redis_key, time_window 2) 4. 统计当前窗口内的行为次数 (ZCARD) count = self.redis_client.zcard(redis_key) logger.info(f"用户 {user_id} 事件 {event_type} 当前计数: {count}, 阈值: {threshold}") 5. 判断是否触发激活 if count >= threshold: 触发后可选择重置计数或保留,这里选择保留以记录热度,但通过逻辑控制不重复触发 实际生产中可增加“已触发”标记位,防止同一次热度的重复通知 return True, action_name return False, None def execute_action(self, user_id, action_name): """ 执行激活后的具体动作 """ logger.info(f"!!! 触发激活 !!! 用户: {user_id}, 动作: {action_name}") 在此处对接邮件API、短信网关或CRM系统 示例:调用第三方发送优惠券 if action_name == 'send_vip_coupon': self._send_coupon(user_id) elif action_name == 'push_sales_notification': self._notify_sales(user_id) def _send_coupon(self, user_id): 模拟发送优惠券逻辑 logger.info(f"[API调用] 向用户 {user_id} 发送VIP优惠券成功") def _notify_sales(self, user_id): 模拟通知销售逻辑 logger.info(f[API调用] 通知销售团队,用户 {user_id} 有高意向") ```现在我们需要一个 HTTP 接口来接收前端或客户端上报的用户行为数据。这里使用 Flask 构建一个轻量级服务。
```python app.py from flask import Flask, request, jsonify from activation_engine import ActivationEngine app = Flask(__name__) engine = ActivationEngine() @app.route('/track', methods=['POST']) def track_event(): """ 接收行为上报接口 期望 JSON 格式: { "user_id": "user_12345", "event_type": "view_pricing_page" } """ try: data = request.get_json() if not data: return jsonify({"status": "error", "message": "请求数据为空"}), 400 user_id = data.get('user_id') event_type = data.get('event_type') if not user_id or not event_type: return jsonify({"status": "error", "message": "缺少 user_id 或 event_type 参数"}), 400 调用引擎处理 is_triggered, action_name = engine.record_event_and_check(user_id, event_type) response_data = { "status": "success", "user_id": user_id, "event_type": event_type, "triggered": is_triggered } if is_triggered: engine.execute_action(user_id, action_name) response_data["action_executed"] = action_name return jsonify(response_data), 200 except Exception as e: return jsonify({"status": "error", "message": str(e)}), 500 if __name__ == '__main__': 启动服务,监听 5000 端口 app.run(host='0.0.0.0', port=5000, debug=True) ```代码编写完毕,现在我们启动系统并进行模拟测试,验证滑动窗口逻辑是否生效。
1. 启动 Flask 服务
在项目根目录下执行:
```bash python app.py ```看到输出 Running on http://0.0.0.0:5000 表示服务启动成功。
2. 使用 cURL 进行模拟测试
我们模拟一个用户 test_user_01 连续访问定价页面。根据配置,阈值是 3 次,窗口是 60 秒。我们需要发送 3 次请求。
打开一个新的终端窗口,执行以下命令:

第 1 次访问:
```bash curl -X POST http://localhost:5000/track \ -H "Content-Type: application/json" \ -d '{"user_id": "test_user_01", "event_type": "view_pricing_page"}' ```预期返回:"triggered": false。日志显示计数为 1。
第 2 次访问:
```bash curl -X POST http://localhost:5000/track \ -H "Content-Type: application/json" \ -d '{"user_id": "test_user_01", "event_type": "view_pricing_page"}' ```预期返回:"triggered": false。日志显示计数为 2。
第 3 次访问(触发激活):
```bash curl -X POST http://localhost:5000/track \ -H "Content-Type: application/json" \ -d '{"user_id": "test_user_01", "event_type": "view_pricing_page"}' ```预期返回:"triggered": true,且包含 "action_executed": "send_vip_coupon"。
观察 Flask 后台日志,应该能看到 !!! 触发激活 !!! 以及 [API调用] 向用户 test_user_01 发送VIP优惠券成功 的字样。
3. 验证滑动窗口时效性
为了验证系统不会无限期累计次数,我们可以修改 config.py 中的 time_window 为 5 秒(为了快速测试)。重启服务。
连续发送 3 次请求触发激活。等待 6 秒后,再次发送 1 次请求。你会发现计数重置为 1,且不会触发激活。这说明 Redis 的 ZREMRANGEBYSCORE 命令正确清理了过期数据。
以上代码完成了核心逻辑闭环,但在生产环境中部署时,还需注意以下技术细节:
1. Redis 持久化配置
编辑 redis.conf,开启 AOF 持久化以确保数据安全,防止重启丢失行为记录:
2. 接口限流
在 app.py 中建议增加 Flask-Limiter 扩展,防止恶意用户高频刷接口导致系统崩溃。安装依赖:
并在代码中初始化:
```python from flask_limiter import Limiter from flask_limiter.util import get_remote_address limiter = Limiter(app, key_func=get_remote_address) 在 track_event 装饰器添加: @limiter.limit("200 per minute") ```3. 异步执行动作
当前 execute_action 是同步执行的,如果对接第三方短信接口较慢,会阻塞请求响应。建议引入 Celery 或 Python 的 threading 将动作执行异步化。
简单的异步改造示例(修改 execute_action 调用处):
至此,一个基于 Python + Redis 的精准流量激活系统已完全构建完成。你可以直接将此代码集成到现有的数据采集 SDK 或后端服务中,实现毫秒级的用户意向判断与实时激活。
易频IT社区是综合性互联网IT技术门户网站,专注分享网络技术、服务器运维、网络安全、编程开发、系统架构、云计算、大数据等行业干货,实时更新IT行业资讯、零基础教程、实战案例,为IT从业者、技术爱好者提供专业的学习交流平台。
Copyright © 2021-2026 易频IT社区. All Rights Reserved. 备案号:闽ICP备2023013482号 网站地图