游戏数据分析与反作弊
1. 游戏数据分析体系
1.1 数据采集
数据采集是游戏数据分析的基石,决定了分析的上限。游戏数据采集涵盖客户端日志、服务端日志、自定义事件、页面埋点等多种方式。
1.1.1 采集方式
客户端日志
客户端日志采集玩家的设备信息、操作行为、性能数据和异常信息。通过 SDK 嵌入游戏客户端,将日志上报至服务端。
// 客户端日志采集示例
public class ClientLogger {
private static final String TAG = "GameAnalytics";
private LogQueue logQueue = new LogQueue(1024); // 内存队列,批量上报
public void logEvent(String eventName, Map<String, Object> params) {
LogEntry entry = new LogEntry();
entry.setEventName(eventName);
entry.setParams(params);
entry.setTimestamp(System.currentTimeMillis());
entry.setUserId(UserManager.getInstance().getUserId());
entry.setSessionId(SessionManager.getCurrentSessionId());
entry.setDeviceId(DeviceInfoUtil.getDeviceId());
entry.setLevel(UserManager.getInstance().getPlayerLevel());
logQueue.enqueue(entry);
// 本地缓存,防止数据丢失
LocalCache.save(entry);
}
// 关键事件打点
public void onLogin() {
Map<String, Object> params = new HashMap<>();
params.put("login_type", getLoginType());
params.put("is_new_device", isNewDevice());
params.put("channel", ChannelUtil.getChannel());
logEvent("login", params);
}
public void onPurchase(String itemId, int amount, String currency, double price) {
Map<String, Object> params = new HashMap<>();
params.put("item_id", itemId);
params.put("amount", amount);
params.put("currency", currency);
params.put("price", price);
params.put("total_spent", UserManager.getInstance().getTotalRecharge());
logEvent("purchase", params);
}
public void onLevelUp(int newLevel) {
Map<String, Object> params = new HashMap<>();
params.put("new_level", newLevel);
params.put("total_play_time", SessionManager.getTotalPlayTime());
params.put("main_quest_progress", QuestManager.getMainQuestProgress());
logEvent("level_up", params);
}
}服务端日志
服务端日志由游戏服务器直接产生,记录玩家在服务器上的所有操作。相比客户端日志,服务端日志更可靠,不可被玩家篡改。
// 服务端日志记录示例
@Component
public class ServerLogCollector {
private final Logger logger = LoggerFactory.getLogger("GameServerLog");
private final KafkaTemplate<String, String> kafkaTemplate;
public void recordPlayerAction(PlayerAction action) {
ServerLogEntry entry = ServerLogEntry.builder()
.eventId(UUID.randomUUID().toString())
.eventType(action.getType())
.userId(action.getUserId())
.roleId(action.getRoleId())
.serverId(action.getServerId())
.timestamp(System.currentTimeMillis())
.ip(action.getClientIp())
.actionDetail(JSON.toJSONString(action.getDetail()))
.build();
// 写入本地日志文件
logger.info(JSON.toJSONString(entry));
// 发送到消息队列
kafkaTemplate.send("game-server-log", entry.getUserId(), JSON.toJSONString(entry));
}
}自定义事件
针对特定的业务场景定义事件,例如关卡开始、关卡结束、BOSS 击杀、装备强化等。
# 自定义事件定义示例
events:
quest_start:
description: "任务开始"
properties:
- name: quest_id
type: string
- name: quest_type
type: string # main/side/daily/activity
- name: player_level
type: int
quest_complete:
description: "任务完成"
properties:
- name: quest_id
type: string
- name: duration_seconds
type: int
- name: rewards
type: array<object>
properties:
- name: item_id
type: string
- name: quantity
type: int
boss_kill:
description: "BOSS击杀"
properties:
- name: boss_id
type: string
- name: boss_level
type: int
- name: party_size
type: int
- name: kill_duration
type: int
- name: damage_dealt
type: long
- name: is_first_kill
type: boolean页面埋点 / 全埋点 / 可视化埋点
- 页面埋点:在游戏特定页面(商城、活动页、充值页)手动埋入代码,追踪 PV/UV、点击热力、转化路径。
- 全埋点:通过 SDK 自动采集玩家所有交互事件,无需手动编码。适用于探索性分析阶段。
- 可视化埋点:通过运营后台界面圈选页面元素,系统自动生成埋点代码。降低埋点成本,提升运营效率。
// 全埋点 SDK 自动采集示例
public class AutoTracker {
// 自动采集点击事件
@Autowired
private TrackerService trackerService;
@Around("@annotation(com.game.sdk.TrackClick)")
public Object trackClick(ProceedingJoinPoint pjp) throws Throwable {
Signature sig = pjp.getSignature();
MethodSignature msig = (MethodSignature) sig;
TrackClick annotation = msig.getMethod().getAnnotation(TrackClick.class);
Object result = pjp.proceed();
// 自动采集事件信息
Map<String, Object> props = new HashMap<>();
props.put("element_id", annotation.elementId());
props.put("element_type", annotation.elementType());
props.put("page_name", annotation.pageName());
props.put("click_result", result != null ? "success" : "fail");
trackerService.track("ui_click", props);
return result;
}
}服务器端 SDK
对于非游戏核心战斗场景(如官网、社区、客服系统),使用服务器端 SDK 直接上报数据。
# 服务器端 SDK 数据上报示例
import requests
import json
import time
class ServerSDK:
def __init__(self, app_id, app_secret, endpoint):
self.app_id = app_id
self.app_secret = app_secret
self.endpoint = endpoint
self.batch_events = []
def track(self, event_name, user_id, properties=None):
"""上报事件"""
event = {
"app_id": self.app_id,
"event": event_name,
"user_id": user_id,
"time": int(time.time() * 1000),
"properties": properties or {},
"$device_id": self._get_device_id(user_id)
}
# 加入批量队列
self.batch_events.append(event)
if len(self.batch_events) >= 50:
self.flush()
def flush(self):
"""批量上报"""
if not self.batch_events:
return
data = {
"app_id": self.app_id,
"sign": self._sign(),
"events": self.batch_events
}
try:
resp = requests.post(
self.endpoint + "/batch",
json=data,
timeout=3,
headers={"Content-Type": "application/json"}
)
if resp.status_code == 200:
self.batch_events = []
except Exception as e:
# 失败重试,写入本地文件
self._save_to_local(data)1.1.2 日志上报协议
HTTP/gRPC 上报
// gRPC 日志上报协议定义
syntax = "proto3";
package game.analytics;
service LogService {
// 单条上报
rpc ReportLog(LogEntry) returns (LogResponse);
// 批量上报
rpc ReportLogBatch(stream LogEntry) returns (LogBatchResponse);
}
message LogEntry {
string event_id = 1;
string user_id = 2;
string role_id = 3;
int32 server_id = 4;
string event_name = 5;
int64 timestamp = 6;
map<string, string> properties = 7;
string device_id = 8;
string session_id = 9;
string ip = 10;
string user_agent = 11;
string channel = 12;
}
message LogResponse {
int32 code = 1;
string message = 2;
}
message LogBatchResponse {
int32 code = 1;
int32 accepted_count = 2;
repeated string failed_ids = 3;
}批量上报与实时上报
| 特性 | 批量上报 | 实时上报 |
|---|---|---|
| 延迟 | 分钟级 | 秒级 |
| 吞吐量 | 高 | 中 |
| 资源消耗 | 低 | 高 |
| 适用场景 | 行为日志、事件追踪 | 实时监控、反作弊告警 |
// 批量上报实现
@Component
public class BatchReporter {
private final BlockingQueue<LogEntry> queue = new LinkedBlockingQueue<>(10000);
private final ExecutorService executor = Executors.newSingleThreadExecutor();
@PostConstruct
public void init() {
executor.submit(() -> {
List<LogEntry> batch = new ArrayList<>(500);
while (true) {
// 每 5 秒或满 500 条上报
try {
Thread.sleep(5000);
queue.drainTo(batch, 500);
if (!batch.isEmpty()) {
reportBatch(batch);
batch.clear();
}
} catch (Exception e) {
log.error("Batch report failed", e);
}
}
});
}
public void add(LogEntry entry) {
queue.offer(entry);
}
private void reportBatch(List<LogEntry> batch) {
// HTTP 批量上报
httpClient.post("https://analytics.game.com/batch")
.body(batch)
.execute();
}
}采样策略
对于海量日志(如战斗帧数据、位置数据),采用采样策略降低数据量。
public class Sampler {
// 按比例采样
public static boolean sampleByRate(double rate) {
return ThreadLocalRandom.current().nextDouble() < rate;
}
// 按用户 ID 哈希采样
public static boolean sampleByUserHash(String userId, int modulus, int target) {
return Math.abs(userId.hashCode()) % modulus == target;
}
// 分层采样:低等级玩家 10%,高等级玩家 100%
public static boolean stratifiedSample(int playerLevel, double lowLevelRate, int threshold) {
if (playerLevel >= threshold) {
return true; // 高等级全量
}
return sampleByRate(lowLevelRate);
}
// 自适应采样
public static boolean adaptiveSample(String eventName, long currentCount, long threshold) {
if (currentCount < threshold) {
return true; // 数据量小全量
}
// 数据量大后逐步降低采样率
double rate = Math.max(0.01, 1.0 - Math.log10(currentCount - threshold + 1) * 0.1);
return sampleByRate(rate);
}
}1.1.3 数据 ETL
-- 数据清洗示例:去除异常数据
INSERT OVERWRITE TABLE ods_game_event PARTITION(dt='${dt}')
SELECT
event_id,
user_id,
role_id,
server_id,
event_name,
timestamp,
properties,
device_id,
session_id
FROM raw_game_event
WHERE dt = '${dt}'
AND user_id IS NOT NULL
AND user_id != ''
AND LENGTH(user_id) < 64
AND timestamp > 0
AND timestamp < UNIX_TIMESTAMP() * 1000
AND server_id > 0;数据格式统一
# 数据格式统一 ETL 示例
def normalize_event(raw_event):
"""统一不同来源的事件格式"""
normalized = {
"event_id": raw_event.get("event_id") or str(uuid.uuid4()),
"event_name": raw_event.get("event_name", raw_event.get("event", "")),
}
# 统一时间格式(毫秒时间戳)
raw_time = raw_event.get("timestamp", raw_event.get("time", raw_event.get("ts", 0)))
if isinstance(raw_time, str):
normalized["timestamp"] = int(date_parser(raw_time).timestamp() * 1000)
elif raw_time < 1000000000000: # 秒级转毫秒
normalized["timestamp"] = raw_time * 1000
else:
normalized["timestamp"] = raw_time
# 用户 ID 统一
normalized["user_id"] = raw_event.get("user_id") or raw_event.get("uid", "")
normalized["role_id"] = raw_event.get("role_id") or raw_event.get("rid", "")
# 属性字段展开
props = raw_event.get("properties", raw_event.get("params", {}))
if isinstance(props, str):
props = json.loads(props)
normalized["properties"] = props
return normalized1.1.4 用户设备信息
// Android 设备信息采集
public class AndroidDeviceInfoCollector {
public static DeviceInfo collect(Context context) {
DeviceInfo info = new DeviceInfo();
// UUID
info.setUuid(DeviceUUID.getUUID(context));
// 设备 ID
info.setDeviceId(Settings.Secure.getString(context.getContentResolver(),
Settings.Secure.ANDROID_ID));
// 渠道 ID
try {
ApplicationInfo appInfo = context.getPackageManager()
.getApplicationInfo(context.getPackageName(), PackageManager.GET_META_DATA);
info.setChannelId(appInfo.metaData.getString("CHANNEL_ID"));
} catch (Exception e) {
info.setChannelId("unknown");
}
// 广告 ID (GAID)
try {
AdvertisingIdClient.Info adInfo = AdvertisingIdClient.getAdvertisingIdInfo(context);
info.setGaid(adInfo.getId());
info.setLimitAdTracking(adInfo.isLimitAdTrackingEnabled());
} catch (Exception e) {
info.setGaid(null);
}
// OAID (匿名设备标识符)
try {
String oaid = OaidHelper.getOAID(context);
info.setOaid(oaid);
} catch (Exception e) {
info.setOaid(null);
}
// IMEI (需要 READ_PHONE_STATE 权限)
if (PermissionChecker.hasPermission(context, Manifest.permission.READ_PHONE_STATE)) {
TelephonyManager tm = (TelephonyManager) context.getSystemService(Context.TELEPHONY_SERVICE);
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) {
info.setImei(tm.getImei());
} else {
info.setImei(tm.getDeviceId());
}
}
// AndroidID
info.setAndroidId(Settings.Secure.getString(
context.getContentResolver(), Settings.Secure.ANDROID_ID));
return info;
}
}// iOS 设备信息采集 (IDFA)
#import <AdSupport/AdSupport.h>
#import <AppTrackingTransparency/AppTrackingTransparency.h>
- (DeviceInfo *)collectDeviceInfo {
DeviceInfo *info = [[DeviceInfo alloc] init];
// UUID
info.uuid = [[[UIDevice currentDevice] identifierForVendor] UUIDString];
// IDFA
if (@available(iOS 14, *)) {
[ATTrackingManager requestTrackingAuthorizationWithCompletionHandler:^(ATTrackingManagerAuthorizationStatus status) {
if (status == ATTrackingManagerAuthorizationStatusAuthorized) {
info.idfa = [[ASIdentifierManager sharedManager] advertisingIdentifier].UUIDString;
} else {
info.idfa = nil;
}
}];
} else {
if ([[ASIdentifierManager sharedManager] isAdvertisingTrackingEnabled]) {
info.idfa = [[ASIdentifierManager sharedManager] advertisingIdentifier].UUIDString;
}
}
return info;
}隐私合规
// 隐私合规数据采集
public class PrivacyCompliantCollector {
// 隐私策略状态
public enum ConsentStatus {
NOT_ASKED,
GRANTED,
DENIED,
EXPIRED
}
private ConsentStatus analyticsConsent = ConsentStatus.NOT_ASKED;
// 请求用户同意
public void requestConsent(Activity activity) {
PrivacyDialog dialog = new PrivacyDialog(activity);
dialog.setTitle("数据采集授权");
dialog.setMessage("为了提升游戏体验,我们需要采集游戏行为数据。"
+ "我们不会收集您的敏感个人信息。您可以在设置中随时关闭。");
dialog.setPositiveButton("同意", () -> {
analyticsConsent = ConsentStatus.GRANTED;
startAnalytics();
});
dialog.setNegativeButton("拒绝", () -> {
analyticsConsent = ConsentStatus.DENIED;
});
dialog.show();
}
// 根据同意状态决定是否采集
public boolean shouldCollect(String dataType) {
if (analyticsConsent != ConsentStatus.GRANTED) {
return false;
}
// 敏感数据类型单独控制
if ("imei".equals(dataType) || "idfa".equals(dataType)) {
return additionalAdConsent;
}
return true;
}
}1.2 数据仓库分层
1.2.1 分层设计
游戏数据仓库采用标准的分层架构,从原始数据到业务应用逐层加工。
+-------------------------------------------------------+
| ADS 应用层 |
| 玩家价值分层 | 流失预测 | LTV 报表 | 活动分析 |
+-------------------------------------------------------+
| DWS 汇总层 |
| 玩家日汇总 | 道具汇总 | 充值汇总 | 活跃汇总 |
+-------------------------------------------------------+
| DWD 明细层 |
| 事件事实表 | 快照事实表 | 状态事实表 |
+-------------------------------------------------------+
| ODS 操作层 |
| 原始日志 | 服务端日志 | 业务库 Binlog |
+-------------------------------------------------------+ODS(操作数据层)
保留原始数据,不做或少做加工。
-- ODS 层建表示例
CREATE TABLE ods_game_event (
event_id STRING COMMENT '事件ID',
user_id STRING COMMENT '用户ID',
role_id STRING COMMENT '角色ID',
server_id INT COMMENT '服务器ID',
event_name STRING COMMENT '事件名称',
event_time BIGINT COMMENT '事件时间(毫秒)',
properties MAP<STRING, STRING> COMMENT '事件属性',
device_id STRING COMMENT '设备ID',
session_id STRING COMMENT '会话ID',
ip STRING COMMENT 'IP地址',
channel STRING COMMENT '渠道',
client_version STRING COMMENT '客户端版本',
country STRING COMMENT '国家',
province STRING COMMENT '省份',
city STRING COMMENT '城市',
isp STRING COMMENT '运营商',
dt STRING COMMENT '分区日期(yyyyMMdd)'
)
PARTITIONED BY (dt STRING)
STORED AS ORC
TBLPROPERTIES (
"orc.compress" = "SNAPPY",
"orc.bloom.filter.columns" = "user_id,event_name",
"partition.retention.days" = "90"
);DWD(数据明细层)
对 ODS 层数据进行清洗、去重、格式统一、维度退化。
-- DWD 层事实表
CREATE TABLE dwd_event_fact (
event_id STRING COMMENT '事件ID',
user_id STRING COMMENT '用户ID',
role_id STRING COMMENT '角色ID',
server_id INT COMMENT '服务器ID',
event_name STRING COMMENT '事件名称',
event_time BIGINT COMMENT '事件时间(毫秒)',
event_date STRING COMMENT '事件日期(yyyyMMdd)',
event_hour INT COMMENT '事件小时',
-- 退化维度
player_level INT COMMENT '玩家等级',
vip_level INT COMMENT 'VIP等级',
platform STRING COMMENT '平台',
channel STRING COMMENT '渠道',
-- 事件属性展开
item_id STRING COMMENT '道具ID',
item_quantity INT COMMENT '道具数量',
cost_amount DECIMAL(18,2) COMMENT '消费金额',
currency_type STRING COMMENT '货币类型',
-- 技术字段
device_type STRING COMMENT '设备型号',
os_version STRING COMMENT '操作系统版本',
network_type STRING COMMENT '网络类型',
ip STRING COMMENT 'IP地址',
country STRING COMMENT '国家',
province STRING COMMENT '省份',
dt STRING COMMENT '分区日期'
)
PARTITIONED BY (dt STRING)
STORED AS ORC;DWS(数据汇总层)
按日、周、月粒度的玩家汇总数据。
-- DWS 层玩家日汇总表
CREATE TABLE dws_player_daily_agg (
user_id STRING COMMENT '用户ID',
role_id STRING COMMENT '角色ID',
server_id INT COMMENT '服务器ID',
dt STRING COMMENT '日期',
-- 活跃指标
login_count INT COMMENT '登录次数',
total_online_sec BIGINT COMMENT '在线时长(秒)',
session_count INT COMMENT '会话次数',
-- 游戏进度
current_level INT COMMENT '当前等级',
exp_gained BIGINT COMMENT '获得经验',
quest_completed INT COMMENT '完成任务数',
dungeon_cleared INT COMMENT '通关副本数',
-- 消费指标
recharge_amount DECIMAL(18,2) COMMENT '充值金额',
recharge_count INT COMMENT '充值次数',
consumption_amount DECIMAL(18,2) COMMENT '消耗金额',
-- 社交指标
friend_count INT COMMENT '好友数',
guild_activity INT COMMENT '公会活跃度',
pvp_battle_count INT COMMENT 'PVP战斗次数',
-- 资源变化
gold_balance BIGINT COMMENT '金币余额',
diamond_balance BIGINT COMMENT '钻石余额',
gold_gained BIGINT COMMENT '获得金币',
gold_spent BIGINT COMMENT '消耗金币',
diamond_gained BIGINT COMMENT '获得钻石',
diamond_spent BIGINT COMMENT '消耗钻石',
-- 设备信息
device_id STRING COMMENT '设备ID',
platform STRING COMMENT '平台',
channel STRING COMMENT '渠道'
)
PARTITIONED BY (dt STRING)
STORED AS ORC;ADS(应用数据层)
面向业务应用的报表数据。
-- ADS 层 LTV 报表
CREATE TABLE ads_ltv_report (
install_date STRING COMMENT '安装日期',
channel STRING COMMENT '渠道',
platform STRING COMMENT '平台',
country STRING COMMENT '国家',
new_user_count BIGINT COMMENT '新增用户数',
-- LTV 累计
ltv_d1 DECIMAL(18,4) COMMENT '第1天LTV',
ltv_d3 DECIMAL(18,4) COMMENT '第3天LTV',
ltv_d7 DECIMAL(18,4) COMMENT '第7天LTV',
ltv_d14 DECIMAL(18,4) COMMENT '第14天LTV',
ltv_d30 DECIMAL(18,4) COMMENT '第30天LTV',
-- 留存率
retention_d1 DECIMAL(5,4) COMMENT '次日留存',
retention_d3 DECIMAL(5,4) COMMENT '3日留存',
retention_d7 DECIMAL(5,4) COMMENT '7日留存',
retention_d14 DECIMAL(5,4) COMMENT '14日留存',
retention_d30 DECIMAL(5,4) COMMENT '30日留存',
-- ARPU
arpu_d7 DECIMAL(18,4) COMMENT '7日ARPU',
arpu_d30 DECIMAL(18,4) COMMENT '30日ARPU',
-- ARPPU
arppu_d7 DECIMAL(18,4) COMMENT '7日ARPPU',
arppu_d30 DECIMAL(18,4) COMMENT '30日ARPPU',
dt STRING COMMENT '统计日期'
)
PARTITIONED BY (dt STRING)
STORED AS ORC;1.2.2 宽表设计
宽表将多个维度和指标合并到一张表中,减少 JOIN 操作,提升查询性能。
-- 玩家日宽表
CREATE TABLE dws_player_daily_wide (
-- 维度字段
user_id STRING,
role_id STRING,
server_id INT,
dt STRING,
platform STRING,
channel STRING,
install_date STRING,
first_recharge_date STRING,
-- 玩家属性
player_level INT,
vip_level INT,
total_recharge DECIMAL(18,2),
total_consumption DECIMAL(18,2),
register_days INT,
-- 当日行为
login_count INT,
online_seconds BIGINT,
pvp_battles INT,
dungeons_cleared INT,
quests_completed INT,
-- 当日消费
recharge_amount DECIMAL(18,2),
recharge_count INT,
gold_income BIGINT,
gold_expense BIGINT,
diamond_income BIGINT,
diamond_expense BIGINT,
-- 累计指标
lifetime_recharge DECIMAL(18,2),
lifetime_consumption DECIMAL(18,2),
lifetime_login_days INT,
last_active_days_ago INT, -- 上次活跃距今天数
-- 标签字段
player_segment STRING COMMENT '玩家分层标签',
churn_risk STRING COMMENT '流失风险',
payment_intent STRING COMMENT '付费意愿',
-- 技术字段
device_id STRING,
ip_address STRING,
country STRING
)
PARTITIONED BY (dt STRING)
STORED AS ORC;1.2.3 分区策略与数据生命周期
-- 分区策略示例
-- 按日期分区,每天一个分区
PARTITIONED BY (dt STRING)
-- 多级分区(减少小文件数)
PARTITIONED BY (year STRING, month STRING, day STRING)
-- 按事件名分区(适合事件表的快速过滤)
PARTITIONED BY (dt STRING, event_name STRING)# 数据生命周期配置
data_lifecycle:
ods_log:
raw_logs: 7 days # 原始日志保留7天
compressed_logs: 30 days # 压缩后保留30天
dwd_fact:
event_fact: 90 days # 事件事实表90天
snapshot_fact: 180 days # 快照事实表180天
dws_agg:
daily_agg: 365 days # 日汇总保留1年
weekly_agg: 2 years # 周汇总保留2年
monthly_agg: 5 years # 月汇总保留5年
ads_report:
permanent: true # 业务报表永久保留
cold_storage:
enable: true
archive_after_days: 180 # 180天后转为冷存储
cold_storage_engine: "OSS/Archive"冷热分离
-- 冷热数据分离查询
-- 热数据(最近30天)使用高性能存储
SELECT * FROM dwd_event_fact
WHERE dt >= DATE_FORMAT(DATE_SUB(CURRENT_DATE, 30), 'yyyyMMdd')
AND user_id = 'target_user';
-- 冷数据(超过30天)使用归档存储
SELECT * FROM dwd_event_fact_archive
WHERE dt < DATE_FORMAT(DATE_SUB(CURRENT_DATE, 30), 'yyyyMMdd')
AND user_id = 'target_user';1.2.4 维度表与事实表
玩家维表
CREATE TABLE dim_player (
user_id STRING COMMENT '用户ID',
role_id STRING COMMENT '角色ID',
role_name STRING COMMENT '角色名',
server_id INT COMMENT '服务器ID',
-- 注册信息
register_time BIGINT COMMENT '注册时间',
register_date STRING COMMENT '注册日期',
register_channel STRING COMMENT '注册渠道',
register_ip STRING COMMENT '注册IP',
register_device STRING COMMENT '注册设备ID',
-- 当前状态
current_level INT COMMENT '当前等级',
current_vip INT COMMENT '当前VIP等级',
total_recharge DECIMAL(18,2) COMMENT '累计充值',
total_consumption DECIMAL(18,2) COMMENT '累计消费',
last_login_time BIGINT COMMENT '最后登录时间',
last_logout_time BIGINT COMMENT '最后登出时间',
-- 属性标签
gender INT COMMENT '性别',
age_group STRING COMMENT '年龄段',
country STRING COMMENT '国家',
city STRING COMMENT '城市',
-- 玩家分层
player_segment STRING COMMENT '玩家分层',
lifecycle_stage STRING COMMENT '生命周期阶段',
payment_tier STRING COMMENT '付费档次',
-- 元数据
etl_time TIMESTAMP COMMENT 'ETL时间',
version INT COMMENT '版本号'
)
STORED AS ORC
TBLPROPERTIES ("transactional" = "true"); -- 支持缓慢变化维更新道具维表
CREATE TABLE dim_item (
item_id STRING COMMENT '道具ID',
item_name STRING COMMENT '道具名称',
item_type STRING COMMENT '道具类型: equipment/consumable/material/card',
item_quality INT COMMENT '品质: 1-白色 2-绿色 3-蓝色 4-紫色 5-橙色',
item_subtype STRING COMMENT '子类型: weapon/armor/accessory/potion',
-- 经济属性
base_price DECIMAL(18,2) COMMENT '基准价格',
sell_price DECIMAL(18,2) COMMENT '出售价格',
max_stack INT COMMENT '最大堆叠数',
is_tradeable BOOLEAN COMMENT '是否可交易',
is_bind_on_pickup BOOLEAN COMMENT '是否拾取绑定',
-- 道具来源
obtain_source ARRAY<STRING> COMMENT '获取来源: [shop,drop,craft,quest]',
drop_rate DECIMAL(10,6) COMMENT '掉落概率',
-- 分类属性
category_path STRING COMMENT '分类路径',
tags ARRAY<STRING> COMMENT '标签',
etl_time TIMESTAMP
)
STORED AS ORC;活动维表
CREATE TABLE dim_activity (
activity_id STRING COMMENT '活动ID',
activity_name STRING COMMENT '活动名称',
activity_type STRING COMMENT '活动类型: recharge/consumption/boss/gvg',
start_time BIGINT COMMENT '开始时间',
end_time BIGINT COMMENT '结束时间',
server_ids ARRAY<INT> COMMENT '参与服务器列表',
-- 活动规则
entry_condition STRING COMMENT '参与条件(JSON)',
reward_rules STRING COMMENT '奖励规则(JSON)',
rank_rules STRING COMMENT '排名规则(JSON)',
-- 运营信息
operator STRING COMMENT '运营负责人',
budget DECIMAL(18,2) COMMENT '活动预算',
expected_arppu DECIMAL(18,2) COMMENT '预期ARPPU',
status STRING COMMENT '状态: draft/published/running/ended',
etl_time TIMESTAMP
)
STORED AS ORC;日期维表
CREATE TABLE dim_date (
date_key STRING COMMENT '日期键 yyyyMMdd',
date_value DATE COMMENT '日期',
year INT COMMENT '年份',
quarter INT COMMENT '季度',
month INT COMMENT '月份',
week INT COMMENT '年内第几周',
day_of_year INT COMMENT '年内第几天',
day_of_month INT COMMENT '月内第几天',
day_of_week INT COMMENT '周内第几天(1-7)',
is_weekend BOOLEAN COMMENT '是否周末',
is_holiday BOOLEAN COMMENT '是否节假日',
season STRING COMMENT '季节',
game_week INT COMMENT '游戏运营周(开服第几周)',
server_life_day INT COMMENT '服务器开服天数'
)
STORED AS ORC;事件事实表
CREATE TABLE dwd_event_fact (
event_id STRING,
user_id STRING,
role_id STRING,
server_id INT,
event_name STRING,
-- 时间维度
event_time BIGINT,
dt STRING,
hour INT,
minute INT,
weekday INT,
-- 维度外键
date_key STRING COMMENT '日期维表FK',
player_key STRING COMMENT '玩家维表FK',
item_key STRING COMMENT '道具维表FK',
activity_key STRING COMMENT '活动维表FK',
-- 事件度量
duration_ms BIGINT COMMENT '持续时间',
item_count INT COMMENT '道具数量',
amount DECIMAL(18,2) COMMENT '金额',
count INT COMMENT '计数',
-- 上下文
source STRING COMMENT '事件来源',
target STRING COMMENT '事件目标',
result STRING COMMENT '事件结果',
detail STRING COMMENT '详细信息(JSON)',
dt PARTITION
)
PARTITIONED BY (dt STRING)
STORED AS ORC;累计快照事实表
-- 玩家月度累计快照
CREATE TABLE dws_player_monthly_snapshot (
user_id STRING,
role_id STRING,
server_id INT,
month_key STRING COMMENT '月份 yyyyMM',
-- 累计到当月的指标
cum_login_days INT COMMENT '累计登录天数',
cum_recharge_amount DECIMAL(18,2) COMMENT '累计充值金额',
cum_recharge_count INT COMMENT '累计充值次数',
cum_consumption DECIMAL(18,2) COMMENT '累计消费',
cum_dungeon_cleared INT COMMENT '累计通关副本',
-- 当月指标
month_login_days INT COMMENT '本月登录天数',
month_recharge DECIMAL(18,2) COMMENT '本月充值',
month_online_hours DOUBLE COMMENT '本月在线小时',
-- 月末状态
end_level INT COMMENT '月末等级',
end_vip INT COMMENT '月末VIP',
end_gold_balance BIGINT COMMENT '月末金币余额',
end_diamond_balance BIGINT COMMENT '月末钻石余额',
etl_time TIMESTAMP
)
PARTITIONED BY (month_key STRING)
STORED AS ORC;1.2.5 数据仓库建模方法
星型模型
事实表在中心,维度表呈辐射状围绕。查询性能好,适合 OLAP。
-- 星型模型示例
-- 事实表
CREATE TABLE fact_recharge (
recharge_id STRING,
user_id STRING,
date_key STRING, -- 关联 dim_date
item_key STRING, -- 关联 dim_item
channel_key STRING, -- 关联 dim_channel
amount DECIMAL(18,2),
payment_method STRING,
status STRING
);
-- 维度表
-- dim_date, dim_player, dim_item, dim_channel 分别在一侧雪花模型
维度表进一步规范化,拆分为子维度表。节省存储空间,但查询需要更多 JOIN。
-- 雪花模型示例
-- dim_player 进一步拆分
CREATE TABLE dim_player (
user_id STRING PRIMARY KEY,
register_channel_key STRING, -- 关联 dim_channel
register_device_key STRING, -- 关联 dim_device
location_key STRING, -- 关联 dim_location
...
);
CREATE TABLE dim_device (
device_key STRING PRIMARY KEY,
device_model STRING,
os_type STRING,
os_version STRING,
screen_size STRING,
brand STRING
);
CREATE TABLE dim_location (
location_key STRING PRIMARY KEY,
country STRING,
province STRING,
city STRING,
isp STRING
);事实星座(事实星系)
多个事实表共享维度表,形成星座结构。
-- 事实星座:多个事实表共享维度
-- 共享维度:dim_date, dim_player, dim_server, dim_channel
-- 事实表1:充值事实
CREATE TABLE fact_recharge (
recharge_id STRING,
date_key STRING, -- 共享 dim_date
user_id STRING, -- 共享 dim_player
server_id INT, -- 共享 dim_server
...
);
-- 事实表2:关卡通关事实
CREATE TABLE fact_dungeon (
dungeon_id STRING,
date_key STRING, -- 共享 dim_date
user_id STRING, -- 共享 dim_player
server_id INT, -- 共享 dim_server
...
);
-- 事实表3:社交行为事实
CREATE TABLE fact_social (
social_id STRING,
date_key STRING, -- 共享 dim_date
user_id STRING, -- 共享 dim_player
...
);1.3 核心指标
1.3.1 用户指标
-- DAU 计算
SELECT dt, COUNT(DISTINCT user_id) AS dau
FROM dwd_event_fact
WHERE event_name = 'login'
AND dt >= '20260701'
AND dt <= '20260707'
GROUP BY dt
ORDER BY dt;
-- WAU 计算
SELECT
DATE_FORMAT(DATE_SUB(dt, (DAYOFWEEK(dt) - 1) % 7), 'yyyyMMdd') AS week_start,
COUNT(DISTINCT user_id) AS wau
FROM dwd_event_fact
WHERE event_name = 'login'
GROUP BY DATE_FORMAT(DATE_SUB(dt, (DAYOFWEEK(dt) - 1) % 7), 'yyyyMMdd');
-- MAU 计算
SELECT
DATE_FORMAT(dt, 'yyyyMM') AS month,
COUNT(DISTINCT user_id) AS mau
FROM dwd_event_fact
WHERE event_name = 'login'
GROUP BY DATE_FORMAT(dt, 'yyyyMM');
-- DNU (日新增用户)
SELECT
dt,
COUNT(DISTINCT p.user_id) AS dnu
FROM dim_player p
WHERE p.register_date = dt
AND dt >= '20260701'
GROUP BY dt;留存率
-- 次日/7日/30日留存率
WITH install_cohort AS (
SELECT
user_id,
register_date AS install_date
FROM dim_player
WHERE register_date >= '20260701'
),
login_dates AS (
SELECT DISTINCT
user_id,
dt
FROM dwd_event_fact
WHERE event_name = 'login'
),
retention_calc AS (
SELECT
i.install_date,
COUNT(DISTINCT i.user_id) AS install_users,
COUNT(DISTINCT CASE WHEN DATEDIFF(l.dt, i.install_date) = 1 THEN i.user_id END) AS retention_d1,
COUNT(DISTINCT CASE WHEN DATEDIFF(l.dt, i.install_date) = 6 THEN i.user_id END) AS retention_d7,
COUNT(DISTINCT CASE WHEN DATEDIFF(l.dt, i.install_date) = 29 THEN i.user_id END) AS retention_d30
FROM install_cohort i
LEFT JOIN login_dates l ON i.user_id = l.user_id
GROUP BY i.install_date
)
SELECT
install_date,
install_users,
ROUND(retention_d1 / install_users, 4) AS retention_rate_d1,
ROUND(retention_d7 / install_users, 4) AS retention_rate_d7,
ROUND(retention_d30 / install_users, 4) AS retention_rate_d30
FROM retention_calc
ORDER BY install_date;渠道留存
-- 各渠道留存率对比
SELECT
p.register_channel,
COUNT(DISTINCT p.user_id) AS install_users,
ROUND(COUNT(DISTINCT CASE WHEN DATEDIFF(l.dt, p.register_date) = 1 THEN p.user_id END)
/ COUNT(DISTINCT p.user_id), 4) AS retention_d1,
ROUND(COUNT(DISTINCT CASE WHEN DATEDIFF(l.dt, p.register_date) = 6 THEN p.user_id END)
/ COUNT(DISTINCT p.user_id), 4) AS retention_d7,
ROUND(COUNT(DISTINCT CASE WHEN DATEDIFF(l.dt, p.register_date) = 29 THEN p.user_id END)
/ COUNT(DISTINCT p.user_id), 4) AS retention_d30
FROM dim_player p
LEFT JOIN dwd_event_fact l
ON p.user_id = l.user_id
AND l.event_name = 'login'
AND l.dt >= '20260701'
WHERE p.register_date >= '20260701'
GROUP BY p.register_channel;1.3.2 行为指标
-- 游戏时长分布
SELECT
dt,
CASE
WHEN total_online_sec < 600 THEN '0-10min'
WHEN total_online_sec < 1800 THEN '10-30min'
WHEN total_online_sec < 3600 THEN '30-60min'
WHEN total_online_sec < 7200 THEN '1-2h'
ELSE '2h+'
END AS duration_bucket,
COUNT(DISTINCT user_id) AS player_count
FROM dws_player_daily_agg
WHERE dt = '20260712'
GROUP BY dt,
CASE
WHEN total_online_sec < 600 THEN '0-10min'
WHEN total_online_sec < 1800 THEN '10-30min'
WHEN total_online_sec < 3600 THEN '30-60min'
WHEN total_online_sec < 7200 THEN '1-2h'
ELSE '2h+'
END;
-- 登录频次
SELECT
player_segment,
AVG(login_days_per_week) AS avg_login_days_per_week,
AVG(sessions_per_day) AS avg_sessions_per_day
FROM (
SELECT
user_id,
p.player_segment,
COUNT(DISTINCT dt) / 7 AS login_days_per_week,
AVG(login_count) AS sessions_per_day
FROM dws_player_daily_agg a
JOIN dim_player p ON a.user_id = p.user_id
WHERE a.dt >= DATE_SUK('20260712', 7)
GROUP BY user_id, p.player_segment
) t
GROUP BY player_segment;
-- 关卡通过率
SELECT
dungeon_id,
dungeon_name,
entry_count,
clear_count,
ROUND(clear_count / entry_count, 4) AS pass_rate,
AVG(clear_duration) AS avg_clear_time_sec
FROM (
SELECT
d.dungeon_id,
d.dungeon_name,
COUNT(CASE WHEN e.event_name = 'dungeon_entry' THEN 1 END) AS entry_count,
COUNT(CASE WHEN e.event_name = 'dungeon_clear' THEN 1 END) AS clear_count,
AVG(CASE WHEN e.event_name = 'dungeon_clear' THEN e.duration_ms / 1000 END) AS clear_duration
FROM dwd_event_fact e
JOIN dim_dungeon d ON e.item_id = d.dungeon_id
WHERE e.event_name IN ('dungeon_entry', 'dungeon_clear')
AND e.dt >= '20260701'
GROUP BY d.dungeon_id, d.dungeon_name
) t
ORDER BY pass_rate;
-- 等级分布
SELECT
current_level,
COUNT(DISTINCT user_id) AS player_count,
ROUND(COUNT(DISTINCT user_id) / SUM(COUNT(DISTINCT user_id)) OVER(), 4) AS pct
FROM dws_player_daily_agg
WHERE dt = '20260712'
GROUP BY current_level
ORDER BY current_level;1.3.3 消费指标
ARPU / ARPPU / LTV
-- ARPU (每用户平均收入)
SELECT
dt,
SUM(recharge_amount) / COUNT(DISTINCT user_id) AS arpu
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt;
-- ARPPU (每付费用户平均收入)
SELECT
dt,
SUM(recharge_amount) / COUNT(DISTINCT CASE WHEN recharge_amount > 0 THEN user_id END) AS arppu
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt;
-- LTV (生命周期价值)
WITH cohort AS (
SELECT
user_id,
register_date,
MIN(dt) AS first_recharge_date
FROM dim_player p
LEFT JOIN dwd_event_fact e ON p.user_id = e.user_id AND e.event_name = 'recharge'
WHERE p.register_date >= '20260701'
GROUP BY user_id, register_date
)
SELECT
c.register_date,
COUNT(DISTINCT c.user_id) AS new_users,
-- 第1天 LTV
ROUND(COALESCE(SUM(CASE WHEN DATEDIFF(f.dt, c.register_date) = 0 THEN f.amount END), 0)
/ COUNT(DISTINCT c.user_id), 4) AS ltv_d1,
-- 第7天 LTV
ROUND(COALESCE(SUM(CASE WHEN DATEDIFF(f.dt, c.register_date) BETWEEN 0 AND 6 THEN f.amount END), 0)
/ COUNT(DISTINCT c.user_id), 4) AS ltv_d7,
-- 第30天 LTV
ROUND(COALESCE(SUM(CASE WHEN DATEDIFF(f.dt, c.register_date) BETWEEN 0 AND 29 THEN f.amount END), 0)
/ COUNT(DISTINCT c.user_id), 4) AS ltv_d30
FROM cohort c
LEFT JOIN dwd_event_fact f ON c.user_id = f.user_id AND f.event_name = 'recharge'
GROUP BY c.register_date
ORDER BY c.register_date;LTV 分渠道
-- 渠道 LTV
SELECT
p.register_channel,
p.register_date,
COUNT(DISTINCT p.user_id) AS new_users,
ROUND(COALESCE(SUM(CASE WHEN DATEDIFF(f.dt, p.register_date) BETWEEN 0 AND 6 THEN f.amount END), 0)
/ COUNT(DISTINCT p.user_id), 4) AS ltv_d7,
ROUND(COALESCE(SUM(CASE WHEN DATEDIFF(f.dt, p.register_date) BETWEEN 0 AND 29 THEN f.amount END), 0)
/ COUNT(DISTINCT p.user_id), 4) AS ltv_d30
FROM dim_player p
LEFT JOIN dwd_event_fact f
ON p.user_id = f.user_id
AND f.event_name = 'recharge'
WHERE p.register_date >= '20260701'
GROUP BY p.register_channel, p.register_date;付费率与付费留存
-- 付费率
SELECT
dt,
COUNT(DISTINCT CASE WHEN recharge_amount > 0 THEN user_id END) AS paying_users,
COUNT(DISTINCT user_id) AS active_users,
ROUND(COUNT(DISTINCT CASE WHEN recharge_amount > 0 THEN user_id END)
/ COUNT(DISTINCT user_id), 4) AS payment_rate
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt;
-- 付费留存(首次付费后持续付费情况)
WITH first_payment AS (
SELECT
user_id,
MIN(dt) AS first_pay_date
FROM dwd_event_fact
WHERE event_name = 'recharge'
GROUP BY user_id
)
SELECT
DATEDIFF(f.dt, fp.first_pay_date) AS days_after_first_pay,
COUNT(DISTINCT fp.user_id) AS cohort_users,
COUNT(DISTINCT CASE WHEN f.amount > 0 THEN fp.user_id END) AS retained_payers,
ROUND(COUNT(DISTINCT CASE WHEN f.amount > 0 THEN fp.user_id END)
/ COUNT(DISTINCT fp.user_id), 4) AS payment_retention_rate
FROM first_payment fp
JOIN dws_player_daily_agg f ON fp.user_id = f.user_id
WHERE fp.first_pay_date >= '20260701'
GROUP BY DATEDIFF(f.dt, fp.first_pay_date)
ORDER BY days_after_first_pay;1.3.4 转化漏斗与归因分析
转化漏斗
-- 用户转化漏斗:展示 -> 点击 -> 注册 -> 创角 -> 付费
SELECT
step,
COUNT(DISTINCT user_id) AS users,
ROUND(COUNT(DISTINCT user_id) / MAX(COUNT(DISTINCT user_id)) OVER(), 4) AS conversion_rate,
ROUND(COUNT(DISTINCT user_id) / LAG(COUNT(DISTINCT user_id)) OVER(ORDER BY step_order), 4) AS step_conversion
FROM (
SELECT 'ad_impression' AS step, 1 AS step_order, device_id AS user_id
FROM dwd_event_fact WHERE event_name = 'ad_impression' AND dt = '20260712'
UNION ALL
SELECT 'ad_click', 2, device_id
FROM dwd_event_fact WHERE event_name = 'ad_click' AND dt = '20260712'
UNION ALL
SELECT 'register', 3, user_id
FROM dwd_event_fact WHERE event_name = 'register' AND dt = '20260712'
UNION ALL
SELECT 'create_role', 4, user_id
FROM dwd_event_fact WHERE event_name = 'create_role' AND dt = '20260712'
UNION ALL
SELECT 'first_payment', 5, user_id
FROM dwd_event_fact WHERE event_name = 'first_recharge' AND dt = '20260712'
) t
GROUP BY step, step_order
ORDER BY step_order;归因分析
// 归因分析引擎
public class AttributionEngine {
public enum AttributionModel {
FIRST_TOUCH, // 首次触点归因
LAST_TOUCH, // 末次触点归因
LINEAR, // 线性归因
TIME_DECAY, // 时间衰减归因
U_SHAPE // U型归因(首次+末次各40%,中间20%)
}
// 首次付费归因
public AttributionResult firstPaymentAttribution(String userId, AttributionModel model) {
List<ClickEvent> touchPoints = getTouchPointsBeforeEvent(userId, "first_recharge");
if (touchPoints.isEmpty()) {
return AttributionResult.unattributed(userId);
}
switch (model) {
case FIRST_TOUCH:
return firstTouchAttribution(userId, touchPoints);
case LAST_TOUCH:
return lastTouchAttribution(userId, touchPoints);
case LINEAR:
return linearAttribution(userId, touchPoints);
case TIME_DECAY:
return timeDecayAttribution(userId, touchPoints);
case U_SHAPE:
return uShapeAttribution(userId, touchPoints);
default:
return lastTouchAttribution(userId, touchPoints);
}
}
private AttributionResult firstTouchAttribution(String userId, List<ClickEvent> touchPoints) {
ClickEvent first = touchPoints.stream()
.min(Comparator.comparingLong(ClickEvent::getTimestamp))
.get();
return AttributionResult.builder()
.userId(userId)
.attributedChannel(first.getChannel())
.attributedCampaign(first.getCampaign())
.model("first_touch")
.weights(singletonMap(first.getChannel(), 1.0))
.build();
}
private AttributionResult timeDecayAttribution(String userId, List<ClickEvent> touchPoints) {
long firstTime = touchPoints.stream()
.mapToLong(ClickEvent::getTimestamp).min().orElse(0);
long lastTime = touchPoints.stream()
.mapToLong(ClickEvent::getTimestamp).max().orElse(0);
long totalSpan = Math.max(lastTime - firstTime, 1);
Map<String, Double> weights = new HashMap<>();
touchPoints.forEach(tp -> {
double position = (double)(tp.getTimestamp() - firstTime) / totalSpan;
double weight = Math.exp(position * 2) / touchPoints.size(); // 时间衰减指数权重
weights.merge(tp.getChannel(), weight, Double::sum);
});
// 归一化
double totalWeight = weights.values().stream().mapToDouble(Double::doubleValue).sum();
weights.replaceAll((k, v) -> v / totalWeight);
String topChannel = weights.entrySet().stream()
.max(Map.Entry.comparingByValue()).get().getKey();
return AttributionResult.builder()
.userId(userId)
.attributedChannel(topChannel)
.model("time_decay")
.weights(weights)
.build();
}
}2. 数据产品与工具
2.1 数据平台架构
+------------------------------------------------------------------+
| 数据产品层 |
| 报表系统 | 自助分析 | OLAP多维分析 | 用户画像 | 精准运营 |
+------------------------------------------------------------------+
| 数据计算层 |
| 离线计算 (Hive/Spark) | 实时计算 (Flink) | 即席查询 |
| | OLAP (Doris/ClickHouse) | |
+------------------------------------------------------------------+
| 数据仓库层 |
| ODS | DWD | DWS | ADS | 维度表 | 宽表 |
+------------------------------------------------------------------+
| 数据采集层 |
| 客户端SDK | 服务端日志 | 数据库Binlog | 业务API | 第三方数据 |
+------------------------------------------------------------------+2.1.1 离线计算
# PySpark 离线 ETL 示例
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
spark = SparkSession.builder \
.appName("GameDailyETL") \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.enableHiveSupport() \
.getOrCreate()
# ODS -> DWD 数据清洗
def etl_ods_to_dwd(event_date):
df = spark.table(f"ods.ods_game_event") \
.filter(col("dt") == event_date) \
.filter(col("user_id").isNotNull()) \
.filter(col("user_id") != "") \
.filter(col("event_name").isNotNull()) \
.dropDuplicates(["event_id"])
# 维度退化:关联玩家维度信息
player_dim = spark.table("dim.dim_player") \
.select("user_id", "current_level", "register_channel")
dwd_df = df.join(player_dim, "user_id", "left") \
.withColumn("event_date", col("dt")) \
.withColumn("event_hour", hour(from_unixtime(col("event_time") / 1000))) \
.select(
"event_id", "user_id", "role_id", "server_id",
"event_name", "event_time", "event_date", "event_hour",
"current_level", "register_channel",
"device_id", "ip", "properties"
)
dwd_df.write \
.mode("overwrite") \
.partitionBy("event_date") \
.format("orc") \
.saveAsTable("dwd.dwd_event_fact")
# 每日 DWS 汇总
def etl_dws_daily_agg(event_date):
df = spark.table("dwd.dwd_event_fact") \
.filter(col("event_date") == event_date)
daily_agg = df.groupBy("user_id", "role_id", "server_id") \
.agg(
count(when(col("event_name") == "login", 1)).alias("login_count"),
sum(when(col("event_name") == "recharge", col("properties.amount"))).alias("recharge_amount"),
count(when(col("event_name") == "recharge", 1)).alias("recharge_count"),
count(when(col("event_name") == "dungeon_clear", 1)).alias("dungeon_cleared"),
countDistinct(when(col("event_name") == "login_session", col("session_id"))).alias("session_count")
) \
.withColumn("dt", lit(event_date))
daily_agg.write \
.mode("overwrite") \
.partitionBy("dt") \
.format("orc") \
.saveAsTable("dws.dws_player_daily_agg")2.1.2 实时计算
// Flink 实时计算示例:实时活跃监控
public class RealtimeActiveMonitor {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Kafka 数据源
FlinkKafkaConsumer<String> kafkaSource = new FlinkKafkaConsumer<>(
"game-server-log",
new SimpleStringSchema(),
getKafkaProperties()
);
DataStream<LogEntry> logStream = env
.addSource(kafkaSource)
.map(json -> JSON.parseObject(json, LogEntry.class));
// 5分钟窗口计算活跃用户
logStream
.filter(log -> "login".equals(log.getEventName()))
.map(log -> Tuple2.of(log.getUserId(), 1L))
.returns(Types.TUPLE(Types.STRING, Types.LONG))
.keyBy(0)
.window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
.reduce((a, b) -> Tuple2.of(a.f0, a.f1 + b.f1))
.map(t -> String.format("{\"userId\":\"%s\",\"eventCount\":%d,\"time\":%d}",
t.f0, t.f1, System.currentTimeMillis()))
.addSink(new FlinkKafkaProducer<>(
"realtime-active-users", new SimpleStringSchema(), getKafkaProperties()
));
// 实时充值监控
logStream
.filter(log -> "recharge".equals(log.getEventName()))
.map(log -> {
JSONObject props = JSON.parseObject(log.getProperties());
return Tuple3.of(log.getServerId(),
props.getDouble("amount"), 1L);
})
.returns(Types.TUPLE(Types.INT, Types.DOUBLE, Types.LONG))
.keyBy(0)
.window(TumblingProcessingTimeWindows.of(Time.minutes(1)))
.aggregate(new SumAggregate())
.map(t -> String.format(
"{\"serverId\":%d,\"totalAmount\":%.2f,\"transactionCount\":%d,\"time\":%d}",
t.f0, t.f1, t.f2, System.currentTimeMillis()))
.addSink(new FlinkKafkaProducer<>(
"realtime-recharge", new SimpleStringSchema(), getKafkaProperties()
));
env.execute("RealtimeGameMonitor");
}
public static class SumAggregate
implements AggregateFunction<Tuple3<Integer, Double, Long>,
Tuple3<Integer, Double, Long>,
Tuple3<Integer, Double, Long>> {
@Override
public Tuple3<Integer, Double, Long> createAccumulator() {
return Tuple3.of(0, 0.0, 0L);
}
@Override
public Tuple3<Integer, Double, Long> add(
Tuple3<Integer, Double, Long> value,
Tuple3<Integer, Double, Long> accumulator) {
accumulator.f0 = value.f0;
accumulator.f1 += value.f1;
accumulator.f2 += value.f2;
return accumulator;
}
@Override
public Tuple3<Integer, Double, Long> getResult(
Tuple3<Integer, Double, Long> accumulator) {
return accumulator;
}
@Override
public Tuple3<Integer, Double, Long> merge(
Tuple3<Integer, Double, Long> a,
Tuple3<Integer, Double, Long> b) {
return Tuple3.of(a.f0, a.f1 + b.f1, a.f2 + b.f2);
}
}
}2.1.3 OLAP 引擎对比
| 特性 | Doris | ClickHouse | Presto/Trino | Hive |
|---|---|---|---|---|
| 查询延迟 | 秒级 | 毫秒-秒级 | 秒级-分钟级 | 分钟级 |
| 并发能力 | 高 | 中 | 中 | 低 |
| 数据导入 | 实时导入 | 批量为主 | 无存储 | 批量 |
| SQL兼容性 | MySQL协议 | 自建SQL | ANSI SQL | HiveQL |
| 索引支持 | 前缀索引/布隆过滤器 | 跳数索引 | 依赖数据源 | 分区/桶 |
| 更新能力 | 支持更新 | 不支持 | 不支持 | 支持 |
| 适用场景 | 报表/即席查询 | 实时监控/大宽表 | 跨数据源查询 | 离线批处理 |
-- Doris OLAP 查询示例:即席分析
SELECT
channel,
COUNT(DISTINCT user_id) AS dau,
SUM(recharge_amount) AS total_revenue,
SUM(recharge_amount) / COUNT(DISTINCT user_id) AS arpu
FROM dws_player_daily_agg
WHERE dt >= '20260701'
AND dt <= '20260712'
GROUP BY channel
ORDER BY total_revenue DESC;
-- ClickHouse 实时聚合表
CREATE TABLE game_realtime_agg (
event_time DateTime,
server_id UInt32,
event_name String,
user_id String,
amount Float64
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMMDD(event_time)
ORDER BY (event_time, server_id);
-- 物化视图:每分钟聚合
CREATE MATERIALIZED VIEW game_minute_agg
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMMDD(minute)
ORDER BY (minute, server_id)
AS SELECT
toStartOfMinute(event_time) AS minute,
server_id,
countState() AS event_count,
uniqState(user_id) AS unique_users,
sumState(amount) AS total_amount
FROM game_realtime_agg
GROUP BY minute, server_id;2.2 数据可视化
2.2.1 报表系统
# 使用 Superset API 创建报表
import requests
import json
class SupersetReportManager:
def __init__(self, base_url, access_token):
self.base_url = base_url
self.headers = {
"Authorization": f"Bearer {access_token}",
"Content-Type": "application/json"
}
def create_daily_report(self):
"""创建日活跃报表"""
chart_config = {
"datasource_id": 1,
"datasource_type": "table",
"viz_type": "echarts_timeseries_line",
"params": {
"metrics": [{
"expressionType": "SQL",
"sqlExpression": "COUNT(DISTINCT user_id)",
"label": "DAU"
}],
"time_range": "last 7 days",
"granularity": "day",
"series": "channel"
}
}
resp = requests.post(
f"{self.base_url}/api/v1/chart",
headers=self.headers,
json=chart_config
)
return resp.json()
def create_funnel_chart(self):
"""创建转化漏斗"""
funnel_config = {
"viz_type": "funnel",
"params": {
"metrics": ["COUNT(DISTINCT user_id)"],
"groupby": ["event_name"],
"funnel_steps": [
"ad_impression", "ad_click", "register",
"create_role", "first_recharge"
],
"funnel_order": "descending"
}
}
return funnel_config2.2.2 自助分析与 OLAP 多维分析
-- OLAP 多维钻取分析
-- 上卷:按渠道汇总
SELECT channel, SUM(recharge_amount) as revenue
FROM dws_player_daily_agg
WHERE dt = '20260712'
GROUP BY channel;
-- 下钻:渠道下的具体广告来源
SELECT channel, ad_source, SUM(recharge_amount) as revenue
FROM dws_player_daily_agg
WHERE dt = '20260712'
GROUP BY channel, ad_source;
-- 切片:只看 iOS 平台
SELECT channel, SUM(recharge_amount) as revenue
FROM dws_player_daily_agg
WHERE dt = '20260712' AND platform = 'iOS'
GROUP BY channel;
-- 切块:iOS 且等级大于 30
SELECT channel, SUM(recharge_amount) as revenue
FROM dws_player_daily_agg
WHERE dt = '20260712'
AND platform = 'iOS'
AND player_level > 30
GROUP BY channel;2.2.3 可视化图表类型
漏斗图
import matplotlib.pyplot as plt
import pandas as pd
def plot_conversion_funnel(steps, values):
"""绘制转化漏斗"""
fig, ax = plt.subplots(figsize=(10, 6))
# 计算转化率
rates = [1.0]
for i in range(1, len(values)):
rates.append(values[i] / values[i-1])
# 绘制漏斗
colors = ['#1f77b4', '#2ca02c', '#ff7f0e', '#d62728', '#9467bd']
max_width = 0.8
for i, (step, value, rate, color) in enumerate(zip(steps, values, rates, colors)):
width = max_width * (value / values[0])
left = (1 - width) / 2
ax.barh(i, width, left=left, height=0.5, color=color, alpha=0.8)
label = f"{step}\n{value:,} | {rate*100:.1f}%"
ax.text(0.5, i, label, ha='center', va='center', fontsize=10)
ax.set_yticks(range(len(steps)))
ax.set_yticklabels([])
ax.invert_yaxis()
ax.set_xlim(0, 1)
ax.axis('off')
plt.title('用户转化漏斗', fontsize=14)
plt.tight_layout()
return fig热力图与分布图
def plot_heatmap(data, x_label, y_label, title):
"""绘制热力图"""
import seaborn as sns
fig, ax = plt.subplots(figsize=(12, 8))
sns.heatmap(data, annot=True, fmt='.0f', cmap='YlOrRd',
linewidths=0.5, ax=ax)
ax.set_title(title, fontsize=14)
ax.set_xlabel(x_label, fontsize=12)
ax.set_ylabel(y_label, fontsize=12)
plt.tight_layout()
return fig
def plot_distribution(values, bins=50, title='分布图'):
"""绘制分布图"""
fig, ax = plt.subplots(figsize=(10, 5))
ax.hist(values, bins=bins, alpha=0.7, edgecolor='black')
ax.axvline(np.mean(values), color='red', linestyle='--', label=f'均值: {np.mean(values):.1f}')
ax.axvline(np.median(values), color='green', linestyle='--', label=f'中位数: {np.median(values):.1f}')
ax.set_xlabel('值', fontsize=12)
ax.set_ylabel('频次', fontsize=12)
ax.set_title(title, fontsize=14)
ax.legend()
plt.tight_layout()
return figApache Superset vs Grafana 对比
| 特性 | Apache Superset | Grafana |
|---|---|---|
| 定位 | BI 分析平台 | 监控可视化平台 |
| 数据源 | SQL 数据库、Druid、ClickHouse 等 | Prometheus、Graphite、InfluxDB 等时序库 |
| SQL 支持 | 原生 SQL 查询、SQL Lab | 自建查询语言 |
| 图表类型 | 50+ 种,侧重分析图表 | 30+ 种,侧重时序监控 |
| 权限管理 | 细粒度 RBAC | 基本组织/用户权限 |
| 告警能力 | 有限 | 强,内置告警规则引擎 |
| 适用场景 | 产品运营报表、自助分析 | 系统监控、实时面板 |
2.3 用户画像
2.3.1 标签体系
用户画像标签体系采用分层分类设计,从基础属性到行为偏好逐层细化。
-- 标签表设计
CREATE TABLE dim_user_tag (
user_id STRING COMMENT '用户ID',
tag_id STRING COMMENT '标签ID',
tag_category STRING COMMENT '标签分类: basic/behavior/consumption/social/churn',
tag_name STRING COMMENT '标签名称',
tag_value STRING COMMENT '标签值',
tag_source STRING COMMENT '标签来源: rule/model/manual',
confidence DECIMAL(5,4) COMMENT '置信度',
expire_time BIGINT COMMENT '过期时间',
create_time TIMESTAMP,
update_time TIMESTAMP
)
PARTITIONED BY (dt STRING)
STORED AS ORC;
-- 标签生成规则(SQL 规则引擎)
-- 高价值付费标签
INSERT OVERWRITE TABLE dim_user_tag PARTITION(dt='${dt}', tag_category='consumption')
SELECT
user_id,
'high_value_payer' AS tag_id,
'consumption' AS tag_category,
'高价值付费用户' AS tag_name,
CASE
WHEN lifetime_recharge > 10000 THEN 'S级'
WHEN lifetime_recharge > 5000 THEN 'A级'
WHEN lifetime_recharge > 1000 THEN 'B级'
ELSE 'C级'
END AS tag_value,
'rule' AS tag_source,
1.0 AS confidence,
NULL AS expire_time,
NOW() AS create_time,
NOW() AS update_time
FROM dws_player_daily_wide
WHERE dt = '${dt}'
AND lifetime_recharge > 1000;标签体系分类:
| 标签分类 | 标签示例 | 来源 |
|---|---|---|
| 基础属性 | 性别/年龄/地区/设备/渠道 | 注册信息、设备信息 |
| 行为偏好 | 玩法偏好/PVP/PVE/社交/休闲/时段偏好 | 行为数据分析 |
| 消费特征 | 付费档次/付费频率/付费偏好品类/价格敏感度 | 付费数据分析 |
| 社交特征 | 公会角色/好友数/组队频次/分享意愿 | 社交数据分析 |
| 成长特征 | 升级速度/活跃天数/关卡进度/成就完成率 | 成长数据分析 |
| 风险标签 | 流失风险/作弊嫌疑/退款风险/恶意举报 | 模型预测+规则 |
2.3.2 玩家分层 RFM 模型
-- RFM 模型计算
-- R: 最近一次充值距今天数
-- F: 近30天充值次数
-- M: 近30天充值金额
WITH rfm_raw AS (
SELECT
user_id,
DATEDIFF(CURRENT_DATE, MAX(dt)) AS recency,
COUNT(CASE WHEN recharge_amount > 0 THEN 1 END) AS frequency,
SUM(recharge_amount) AS monetary
FROM dws_player_daily_agg
WHERE dt >= DATE_SUK(CURRENT_DATE, 30)
GROUP BY user_id
),
rfm_score AS (
SELECT
user_id,
recency,
frequency,
monetary,
-- R 分箱:距今天数越小分数越高
CASE
WHEN recency <= 1 THEN 5
WHEN recency <= 3 THEN 4
WHEN recency <= 7 THEN 3
WHEN recency <= 14 THEN 2
ELSE 1
END AS r_score,
-- F 分箱
CASE
WHEN frequency >= 10 THEN 5
WHEN frequency >= 5 THEN 4
WHEN frequency >= 3 THEN 3
WHEN frequency >= 1 THEN 2
ELSE 1
END AS f_score,
-- M 分箱
CASE
WHEN monetary >= 5000 THEN 5
WHEN monetary >= 1000 THEN 4
WHEN monetary >= 500 THEN 3
WHEN monetary >= 100 THEN 2
ELSE 1
END AS m_score
FROM rfm_raw
)
SELECT
user_id,
CONCAT(r_score, f_score, m_score) AS rfm_segment,
r_score, f_score, m_score,
CASE
WHEN r_score >= 4 AND f_score >= 4 AND m_score >= 4 THEN '重要价值用户'
WHEN r_score >= 4 AND f_score >= 3 AND m_score >= 3 THEN '重要发展用户'
WHEN r_score >= 3 AND f_score >= 4 AND m_score >= 4 THEN '重要保持用户'
WHEN r_score <= 2 AND f_score >= 3 AND m_score >= 3 THEN '重要挽留用户'
WHEN r_score >= 4 AND f_score <= 2 THEN '新用户'
WHEN r_score <= 2 AND f_score <= 2 AND m_score <= 2 THEN '流失用户'
ELSE '一般用户'
END AS rfm_label
FROM rfm_score;玩家生命周期管理
玩家生命周期分为五个阶段:导入期、成长期、成熟期、衰退期、流失期。
class PlayerLifecycleManager:
"""玩家生命周期管理"""
STAGES = ['导入期', '成长期', '成熟期', '衰退期', '流失期']
def classify_stage(self, player):
"""判断玩家生命周期阶段"""
register_days = player.register_days
last_active_days = player.last_active_days_ago
level = player.current_level
max_level = self.get_max_level()
# 导入期:注册7天内
if register_days <= 7:
return '导入期'
# 流失期:超过14天未登录
if last_active_days > 14:
return '流失期'
# 衰退期:活跃天数下降,等级增长停滞
if self.is_declining(player):
return '衰退期'
# 成熟期:达到较高等级,持续活跃
if level >= max_level * 0.7:
return '成熟期'
# 成长期:活跃增长中
return '成长期'
def get_stage_action(self, stage):
"""获取各阶段运营策略"""
actions = {
'导入期': {
'goal': '完成新手引导,建立游戏认知',
'actions': ['新手礼包', '引导任务', '首充优惠', '在线奖励'],
'metrics': ['新手完成率', '首日留存', '创角率']
},
'成长期': {
'goal': '加速成长,建立游戏习惯',
'actions': ['成长基金', '等级礼包', '活跃任务', '组队引导'],
'metrics': ['等级提升速度', '日均在线时长', '社交行为']
},
'成熟期': {
'goal': '维持活跃,促进付费',
'actions': ['VIP特惠', '限时活动', '社交竞赛', '新内容推送'],
'metrics': ['ARPU', '付费率', '活动参与率']
},
'衰退期': {
'goal': '召回干预,延缓流失',
'actions': ['回归礼包', '专属活动', '好友召回', '客服关怀'],
'metrics': ['活跃天数变化', '登录频次']
},
'流失期': {
'goal': '流失召回,重新激活',
'actions': ['流失召回邮件', '回归奖励', '版本更新通知', '社交邀请'],
'metrics': ['召回率', '回归后留存']
}
}
return actions.get(stage, {})流失预警与流失干预
# 流失预测模型(伪代码)
class ChurnPredictionModel:
"""基于机器学习的流失预测"""
def __init__(self):
self.model = self._build_model()
self.features = [
'register_days', 'last_active_days_ago',
'avg_online_time_7d', 'login_frequency_7d',
'recharge_amount_7d', 'recharge_count_7d',
'level_growth_7d', 'quest_complete_rate_7d',
'friend_count', 'guild_active_days',
'pvp_battle_count_7d', 'dungeon_clear_count_7d',
'gold_balance_change', 'diamond_balance_change'
]
def _build_model(self):
"""构建梯度提升树模型"""
import lightgbm as lgb
params = {
'objective': 'binary',
'metric': 'auc',
'boosting_type': 'gbdt',
'num_leaves': 63,
'learning_rate': 0.05,
'feature_fraction': 0.8,
'bagging_fraction': 0.8,
'bagging_freq': 5,
'verbose': -1,
'min_data_in_leaf': 100,
'max_depth': 7
}
return lgb.LGBMClassifier(**params)
def predict_churn_risk(self, player_features):
"""
预测流失风险
返回: {user_id: risk_score(0-1)}
risk_score > 0.7: 高危
risk_score 0.4-0.7: 中危
risk_score < 0.4: 低危
"""
risk_score = self.model.predict_proba(player_features)[:, 1]
return risk_score
def get_churn_reason(self, player_features, player_id):
"""使用 SHAP 解释流失原因"""
import shap
explainer = shap.TreeExplainer(self.model)
shap_values = explainer.shap_values(player_features)
# 找出最重要的前3个流失特征
feature_importance = list(zip(self.features, shap_values[0]))
feature_importance.sort(key=lambda x: abs(x[1]), reverse=True)
return {
'player_id': player_id,
'risk_level': 'high' if risk > 0.7 else 'medium' if risk > 0.4 else 'low',
'top_reasons': [
{'feature': f, 'impact': float(s)}
for f, s in feature_importance[:3]
]
}沉默用户召回
-- 沉默用户筛选(7天以上未登录)
SELECT
user_id,
role_name,
current_level,
last_login_time,
DATEDIFF(CURRENT_DATE, last_login_time) AS silent_days,
total_recharge,
player_segment
FROM dim_player
WHERE DATEDIFF(CURRENT_DATE, last_login_time) BETWEEN 7 AND 30
AND DATEDIFF(CURRENT_DATE, last_login_time) <= 30
AND player_segment != 'churned';
-- 召回效果分析
SELECT
recall_channel, -- 召回渠道: SMS/Email/Push/AD
recall_date,
COUNT(DISTINCT user_id) AS recalled_users,
COUNT(DISTINCT CASE WHEN login_days_after_recall >= 1 THEN user_id END) AS day1_active,
COUNT(DISTINCT CASE WHEN login_days_after_recall >= 7 THEN user_id END) AS day7_active,
COUNT(DISTINCT CASE WHEN recharge_amount_after_recall > 0 THEN user_id END) AS repayer_users,
SUM(recharge_amount_after_recall) AS total_recall_revenue
FROM ads_recall_effect
WHERE recall_date >= '20260701'
GROUP BY recall_channel, recall_date;2.3.3 精准运营与 A/B 测试
A/B 测试
// A/B 测试分流器
public class ABTestSplitter {
private static final int BUCKET_COUNT = 10000;
// 确定用户所属实验分组
public static String getExperimentGroup(String userId, String experimentId) {
// 用户哈希一致性:确保同一用户始终进入同一分组
int hash = Math.abs((userId + "_" + experimentId).hashCode());
int bucket = hash % BUCKET_COUNT;
ExperimentConfig config = getExperimentConfig(experimentId);
if (bucket < config.getControlBucketSize()) {
return "control"; // 对照组
} else if (bucket < config.getControlBucketSize() + config.getTreatmentBucketSize()) {
return "treatment"; // 实验组
} else {
return "excluded"; // 不参与实验
}
}
// AA 测试验证
public static boolean validateAATest(String experimentId) {
// 将对照组随机分为两组,验证两组指标无明显差异
List<String> controlUsers = getControlUsers(experimentId);
Map<String, List<String>> aaGroups = splitEvenly(controlUsers);
double groupA_metric = calculateMetric(aaGroups.get("A"));
double groupB_metric = calculateMetric(aaGroups.get("B"));
double pValue = performTTest(groupA_metric, groupB_metric);
// p > 0.05 说明 AA 通过,两组无显著差异
return pValue > 0.05;
}
}
// MAB 多臂老虎机
public class MultiArmedBandit {
// Thompson Sampling 实现
public static String selectVariant(String userId, String experimentId) {
List<Variant> variants = getVariants(experimentId);
Variant bestVariant = null;
double bestSample = -1;
for (Variant v : variants) {
// Beta 分布采样
double sample = BetaDistribution.sample(
v.getSuccessCount() + 1, // alpha: 成功次数 + 1
v.getTrialCount() - v.getSuccessCount() + 1 // beta: 失败次数 + 1
);
if (sample > bestSample) {
bestSample = sample;
bestVariant = v;
}
}
return bestVariant.getId();
}
public static void recordResult(String variantId, boolean converted) {
Variant v = getVariant(variantId);
if (converted) {
v.incrementSuccess();
}
v.incrementTrial();
}
}用户分群与精准运营
-- 用户分群:高付费意愿但近期活跃下降的用户群
INSERT OVERWRITE TABLE ads_user_segment_daily
PARTITION(dt='${dt}', segment_id='high_value_churn_risk')
SELECT
a.user_id,
'high_value_churn_risk' AS segment_id,
CURRENT_TIMESTAMP AS create_time,
JSON_OBJECT(
'reason', '近7天活跃下降但历史付费高',
'last_active_days', a.last_active_days_ago,
'lifetime_recharge', a.lifetime_recharge
) AS segment_meta
FROM dws_player_daily_wide a
WHERE a.dt = '${dt}'
AND a.lifetime_recharge > 1000
AND a.last_active_days_ago BETWEEN 3 AND 7
AND a.online_seconds < 600;3. 反作弊基础
3.1 作弊类型
游戏作弊类型可以分为以下几个大类:
作弊类型
├── 客户端作弊
│ ├── 外挂程序 (Wallhack/Aimbot/透视/自瞄)
│ ├── 脱机挂 (独立的自动脚本客户端)
│ ├── 内存修改 (CE修改器/数值修改)
│ ├── 变速齿轮 (加速/减速/暂停)
│ └── 模拟器/云手机 (模拟操作环境)
├── 硬件辅助
│ ├── 同步器 (一套外设控制多个设备)
│ ├── 连点器/宏 (自动点击/操作录制)
│ ├── 自动脚本 (按键精灵/脚本精灵)
│ └── 物理外挂 (定制硬件设备)
├── 工作室行为
│ ├── 批量建号 (脚本自动注册)
│ ├── 打金/刷资源 (自动化产金)
│ ├── 资源转移 (小号养大号)
│ └── 刷道具 (利用系统漏洞刷取)
└── 业务作弊
├── 恶意退款 (充值后投诉退款)
├── 刷单刷量 (虚假数据刷榜)
├── 虚假流量 (欺诈安装/机刷)
├── 积分墙作弊 (任务奖励欺诈)
├── 刷榜/刷评论 (操纵排行榜)
└── 刷分 (PVP/天梯分数作弊)3.2 检测手段
3.2.1 客户端检测
// 客户端反作弊检测
public class ClientAntiCheat {
private static final Logger log = LoggerFactory.getLogger(ClientAntiCheat.class);
// 1. 内存扫描检测
public static boolean detectMemoryModification() {
// 校验关键内存地址的完整性
long expectedValue = 100;
long actualValue = readMemoryAddress(KEY_ADDRESS_HEALTH);
return actualValue != expectedValue;
}
// 2. 加速检测
public static boolean detectSpeedHack() {
long currentTime = System.currentTimeMillis();
long systemTime = getSystemUptime();
// 检查系统运行时间与游戏内时间的偏差
long gameTimeDelta = currentTime - lastFrameTimestamp;
long systemTimeDelta = systemTime - lastSystemTimestamp;
double ratio = (double) gameTimeDelta / systemTimeDelta;
// 正常范围 0.9 - 1.1,超出则判定加速
return ratio < 0.9 || ratio > 1.1;
}
// 3. 多开检测
public static boolean detectMultiInstance() {
int processCount = 0;
String currentProcess = ProcessUtil.getCurrentProcessName();
for (ProcessInfo proc : ProcessUtil.getRunningProcesses()) {
if (proc.getName().equals(currentProcess)) {
processCount++;
}
}
return processCount > 1;
}
// 4. 模拟器检测
public static boolean detectEmulator() {
// 检测常见模拟器特征
String[] emulatorFiles = {
"/system/bin/ttvd", // 天天模拟器
"/system/bin/wifi", // 夜神模拟器
"/system/lib/libdvm.so", // 蓝叠模拟器
"/system/lib/libnoxspeed.so", // 夜神模拟器
"/data/data/com.mumu.launcher", // 网易MUMU
};
for (String file : emulatorFiles) {
if (new File(file).exists()) {
return true;
}
}
// 检测模拟器特有属性
if (Build.BRAND.contains("genymotion") ||
Build.MANUFACTURER.contains("genymotion") ||
Build.DEVICE.contains("vbox")) {
return true;
}
return false;
}
// 5. Root/越狱检测
public static boolean detectRoot() {
String[] rootPaths = {
"/system/app/Superuser.apk",
"/sbin/su",
"/system/bin/su",
"/system/xbin/su",
"/data/local/xbin/su",
"/data/local/bin/su",
"/system/sd/xbin/su",
"/system/bin/failsafe/su",
"/data/local/su"
};
for (String path : rootPaths) {
if (new File(path).exists()) {
return true;
}
}
return false;
}
// 6. Xposed/Frida 检测
public static boolean detectHookFramework() {
// Xposed 检测
try {
ClassLoader classLoader = ClassLoader.getSystemClassLoader();
Class<?> xposedClass = classLoader.loadClass("de.robv.android.xposed.XposedBridge");
if (xposedClass != null) {
return true;
}
} catch (ClassNotFoundException ignored) {}
// Frida 检测
try {
// 检查 frida-server 端口
Process process = Runtime.getRuntime().exec("netstat -an | findstr 27042");
BufferedReader reader = new BufferedReader(
new InputStreamReader(process.getInputStream()));
String line;
while ((line = reader.readLine()) != null) {
if (line.contains("LISTEN")) {
return true;
}
}
} catch (Exception ignored) {}
return false;
}
// 7. SELinux 验证
public static boolean verifySELinux() {
try {
Process process = Runtime.getRuntime().exec("getenforce");
BufferedReader reader = new BufferedReader(
new InputStreamReader(process.getInputStream()));
String status = reader.readLine();
// 如果 SELinux 被关闭(Permissive 或 Disabled),可能是作弊环境
return "Enforcing".equals(status);
} catch (Exception e) {
return false;
}
}
}3.2.2 服务器端检测
// 服务端反作弊检测引擎
@Component
public class ServerAntiCheatEngine {
private final RuleEngine ruleEngine;
private final BehaviorAnalyzer behaviorAnalyzer;
private final DeviceFingerprintService deviceFingerprintService;
// 1. 行为分析
public DetectionResult analyzeBehavior(PlayerAction action) {
List<String> flags = new ArrayList<>();
// 频率控制:检查操作频率是否异常
if (isAbnormalFrequency(action)) {
flags.add("ABNORMAL_FREQUENCY");
}
// 数据校验:检查请求数据是否合法
if (!validateRequestData(action)) {
flags.add("INVALID_DATA");
}
// 偏差检测:与正常玩家行为对比
if (isBehaviorDeviation(action)) {
flags.add("BEHAVIOR_DEVIATION");
}
return new DetectionResult(action.getUserId(), flags, calculateRiskScore(flags));
}
// 频率控制
private boolean isAbnormalFrequency(PlayerAction action) {
String key = action.getUserId() + ":" + action.getActionType();
long count = redisTemplate.opsForValue().increment(key, 1);
redisTemplate.expire(key, Duration.ofSeconds(1));
Map<String, Integer> thresholds = Map.of(
"attack", 10, // 每秒攻击次数上限
"move", 50, // 每秒移动指令上限
"chat", 5, // 每秒聊天次数上限
"trade", 3 // 每秒交易操作上限
);
int threshold = thresholds.getOrDefault(action.getActionType(), 20);
return count > threshold;
}
// 异常 IP 检测
public boolean isSuspiciousIP(String ip) {
// 检查 IP 是否在黑名单中
if (ipBlacklist.contains(ip)) return true;
// 检查同一 IP 关联的账号数
long accountCount = redisTemplate.opsForHyperLogLog().size("ip:" + ip);
return accountCount > 20; // 同一 IP 超过 20 个账号
}
// 设备指纹检测
public boolean isSuspiciousDevice(String deviceFingerprint) {
// 公用设备检测
long accountCount = deviceFingerprintService.getAccountCount(deviceFingerprint);
if (accountCount > 10) return true;
// 模拟器指纹检测
if (deviceFingerprintService.isEmulator(deviceFingerprint)) return true;
return false;
}
// 黑卡检测
public boolean isBlackCard(String cardNumber) {
// 查询黑卡库
return blackCardCache.containsKey(cardNumber);
}
// 支付作弊 / 退款检测
public void checkPaymentFraud(PaymentEvent payment) {
// 同一账户频繁退款
long refundCount = paymentHistoryDao.countRefunds(
payment.getUserId(),
LocalDate.now().minusDays(30)
);
if (refundCount > 3) {
riskScoreManager.addScore(payment.getUserId(), 50, "频繁退款");
flagForReview(payment.getUserId(), "PAYMENT_FRAUD");
}
// 退款金额占比过高
double refundRate = paymentHistoryDao.getRefundRate(payment.getUserId());
if (refundRate > 0.3) {
flagForReview(payment.getUserId(), "HIGH_REFUND_RATE");
}
}
}3.3 反作弊架构
3.3.1 整体架构
+----------------------------------------------------------+
| 反作弊数据平台 |
| 客户端SDK | 安全服务端 | 数据平台 | 策略引擎 | 人工审核 |
+----------------------------------------------------------+
| 策略引擎 |
| 规则引擎 (Drools) | 机器学习模型 | 实时计算 | 离线分析 |
+----------------------------------------------------------+
| 数据采集与处理 |
| 客户端安全SDK上报 | 服务端行为日志 | 设备指纹 | IP库 |
+----------------------------------------------------------+
| 配置与规则管理 |
| 热更新规则 | 配置下放 | 黑白名单 | 分级策略 | 策略编排 |
+----------------------------------------------------------+3.3.2 策略引擎
// 规则引擎示例 (Drools)
// rules/anti-cheat-rules.drl
package rules.anti_cheat;
import com.game.anticheat.PlayerAction;
import com.game.anticheat.DetectionResult;
rule "同一IP多账号"
when
$ip: String()
$accounts: List(size() > 10) from
accumulate(PlayerAction(ip == $ip, distinctUserId: userId),
collect(distinctUserId))
then
insert(new DetectionResult("SUSPICIOUS_IP_CLUSTER",
"同一IP关联超过10个账号", 60));
end
rule "异常游戏时长"
when
$action: PlayerAction(dailyOnlineMinutes > 1440) // 超过24小时
then
insert(new DetectionResult($action.getUserId(),
"ABNORMAL_ONLINE_TIME", 80));
end
rule "异常充值频率"
when
$action: PlayerAction(rechargeCountPerMinute > 5)
then
insert(new DetectionResult($action.getUserId(),
"ABNORMAL_RECHARGE_FREQUENCY", 70));
end
rule "角色等级与充值不匹配"
when
$action: PlayerAction(level < 10 && totalRecharge > 10000)
then
insert(new DetectionResult($action.getUserId(),
"LEVEL_RECHARGE_MISMATCH", 50));
end
// 热更新规则管理
@Component
public class RuleUpdateManager {
// 从配置中心加载规则
public void refreshRules() {
String rulesYaml = configCenter.getConfig("anticheat.rules");
List<AntiCheatRule> rules = parseRules(rulesYaml);
// 更新规则缓存
ruleCache.putAll(rules);
// 更新 Drools 规则库
kieBase.addPackages(loadRulesFromYaml(rulesYaml));
log.info("反作弊规则已更新,共 {} 条规则", rules.size());
}
// 配置下放
public void deployConfig(String version) {
AntiCheatConfig config = configCenter.getConfig("anticheat.config." + version);
// 下发给客户端
clientConfigSender.send(config.getClientConfig());
// 下发给服务端
serverConfigHolder.update(config.getServerConfig());
// 记录版本号
currentVersion = version;
}
}3.3.3 黑白名单与分级策略
@Component
public class RiskLevelManager {
public enum RiskLevel {
GREEN(0, "正常", "不做处理"),
YELLOW(1, "可疑", "增加监控频率"),
ORANGE(2, "高危", "限制部分功能"),
RED(3, "确定作弊", "封禁处理");
final int level;
final String label;
final String action;
RiskLevel(int level, String label, String action) {
this.level = level;
this.label = label;
this.action = action;
}
}
// 风险评分 -> 等级映射
public RiskLevel classifyRisk(int riskScore) {
if (riskScore >= 80) return RiskLevel.RED;
if (riskScore >= 50) return RiskLevel.ORANGE;
if (riskScore >= 20) return RiskLevel.YELLOW;
return RiskLevel.GREEN;
}
// 分级处置
public void handleRisk(String userId, RiskLevel level, List<String> flags) {
switch (level) {
case YELLOW:
// 增加监控频率,记录详细日志
monitoringService.increaseMonitorLevel(userId);
antiCheatLogger.logSuspicion(userId, flags);
break;
case ORANGE:
// 限制交易、发言等功能
featureRestrictionService.restrict(userId, List.of(
FeatureType.TRADE, FeatureType.CHAT, FeatureType.MAIL
));
// 要求客户端上传更多安全信息
clientSecurityRequester.requestExtraInfo(userId);
antiCheatLogger.logWarning(userId, flags);
break;
case RED:
// 封禁账号
banService.banAccount(userId, flags);
// 设备加入黑名单
deviceBlacklistService.addDevice(userId);
// 通知运营
alertService.sendAlert("封禁通知", userId, flags);
antiCheatLogger.logBan(userId, flags);
break;
}
}
// 人工审核队列
public void submitForReview(String userId, int riskScore, List<String> flags) {
ReviewTicket ticket = ReviewTicket.builder()
.userId(userId)
.riskScore(riskScore)
.flags(flags)
.evidence(collectEvidence(userId))
.timestamp(System.currentTimeMillis())
.priority(riskScore >= 70 ? "high" : "normal")
.build();
reviewQueue.enqueue(ticket);
}
}4. 游戏行为分析
4.1 玩家行为分析
4.1.1 活跃分析
-- 日活跃趋势分析
SELECT
dt,
COUNT(DISTINCT user_id) AS dau,
COUNT(DISTINCT CASE WHEN platform = 'iOS' THEN user_id END) AS ios_dau,
COUNT(DISTINCT CASE WHEN platform = 'Android' THEN user_id END) AS android_dau,
AVG(total_online_sec) AS avg_online_duration,
SUM(total_online_sec) / 3600 AS total_online_hours
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt
ORDER BY dt;
-- 时段活跃分布
SELECT
event_hour,
COUNT(DISTINCT user_id) AS active_users,
COUNT(*) AS event_count
FROM dwd_event_fact
WHERE dt = '20260712'
AND event_name = 'login'
GROUP BY event_hour
ORDER BY event_hour;
-- 在线时长分布
SELECT
CASE
WHEN total_online_sec < 300 THEN '0-5min'
WHEN total_online_sec < 900 THEN '5-15min'
WHEN total_online_sec < 1800 THEN '15-30min'
WHEN total_online_sec < 3600 THEN '30-60min'
WHEN total_online_sec < 7200 THEN '1-2h'
ELSE '2h+'
END AS duration_bucket,
COUNT(DISTINCT user_id) AS player_count,
ROUND(COUNT(DISTINCT user_id) * 100.0 / SUM(COUNT(DISTINCT user_id)) OVER(), 2) AS pct
FROM dws_player_daily_agg
WHERE dt = '20260712'
GROUP BY
CASE
WHEN total_online_sec < 300 THEN '0-5min'
WHEN total_online_sec < 900 THEN '5-15min'
WHEN total_online_sec < 1800 THEN '15-30min'
WHEN total_online_sec < 3600 THEN '30-60min'
WHEN total_online_sec < 7200 THEN '1-2h'
ELSE '2h+'
END
ORDER BY duration_bucket;4.1.2 流失分析
-- 流失前行为特征分析
WITH churned_players AS (
-- 定义流失用户:连续14天未登录
SELECT DISTINCT user_id
FROM dim_player
WHERE last_login_time < UNIX_TIMESTAMP() * 1000 - 14 * 86400 * 1000
AND DATEDIFF(CURRENT_DATE, register_date) > 30
),
behavior_before_churn AS (
SELECT
c.user_id,
DATEDIFF(CAST(FROM_UNIXTIME(MAX(e.event_time / 1000)) AS DATE),
CAST(FROM_UNIXTIME(p.last_login_time / 1000) AS DATE)) AS days_before_churn,
COUNT(CASE WHEN e.event_name = 'login' THEN 1 END) AS login_count,
COUNT(CASE WHEN e.event_name = 'dungeon_clear' THEN 1 END) AS dungeon_clears,
COUNT(CASE WHEN e.event_name = 'pvp_battle' THEN 1 END) AS pvp_battles,
COUNT(CASE WHEN e.event_name = 'recharge' THEN 1 END) AS recharge_count
FROM churned_players c
JOIN dim_player p ON c.user_id = p.user_id
JOIN dwd_event_fact e ON c.user_id = e.user_id
WHERE e.event_time > (p.last_login_time - 7 * 86400 * 1000)
AND e.event_time <= p.last_login_time
GROUP BY c.user_id
)
SELECT
days_before_churn,
AVG(login_count) AS avg_login_count,
AVG(dungeon_clears) AS avg_dungeon_clears,
AVG(pvp_battles) AS avg_pvp_battles,
AVG(recharge_count) AS avg_recharge_count,
COUNT(DISTINCT user_id) AS churned_users
FROM behavior_before_churn
GROUP BY days_before_churn
ORDER BY days_before_churn;活跃天数衰减分析
def analyze_active_day_decay(cohort_data):
"""
分析同期群活跃天数衰减模式
cohort_data: DataFrame with columns [cohort_week, week_index, active_days]
"""
import matplotlib.pyplot as plt
import seaborn as sns
fig, ax = plt.subplots(figsize=(14, 8))
# 每个同期群绘制一条衰减曲线
for cohort in cohort_data['cohort_week'].unique():
data = cohort_data[cohort_data['cohort_week'] == cohort]
ax.plot(data['week_index'], data['active_days'],
marker='o', label=f'{cohort} 同期群', alpha=0.7)
ax.set_xlabel('注册后第N周', fontsize=12)
ax.set_ylabel('周平均活跃天数', fontsize=12)
ax.set_title('同期群活跃天数衰减曲线', fontsize=14)
ax.legend(loc='upper right', fontsize=8, ncol=2)
ax.grid(True, alpha=0.3)
# 添加衰减基准线
from scipy.optimize import curve_fit
def decay_func(x, a, b, c):
return a * np.exp(-b * x) + c
avg_active_days = cohort_data.groupby('week_index')['active_days'].mean()
x_data = np.array(avg_active_days.index)
y_data = np.array(avg_active_days.values)
try:
popt, _ = curve_fit(decay_func, x_data, y_data, maxfev=5000)
ax.plot(x_data, decay_func(x_data, *popt),
'r--', label=f'拟合衰减: y={popt[0]:.1f}*exp(-{popt[1]:.2f}x)+{popt[2]:.1f}',
linewidth=2)
ax.legend()
except:
pass
plt.tight_layout()
return fig关卡通过 / 卡点分析
-- 关卡通过率与卡点识别
SELECT
d.dungeon_id,
d.dungeon_name,
d.difficulty,
d.recommended_level,
COUNT(DISTINCT CASE WHEN e.event_name = 'dungeon_entry' THEN e.user_id END) AS entry_users,
COUNT(DISTINCT CASE WHEN e.event_name = 'dungeon_clear' THEN e.user_id END) AS clear_users,
ROUND(COUNT(DISTINCT CASE WHEN e.event_name = 'dungeon_clear' THEN e.user_id END) * 100.0
/ NULLIF(COUNT(DISTINCT CASE WHEN e.event_name = 'dungeon_entry' THEN e.user_id END), 0), 2) AS pass_rate,
AVG(CASE WHEN e.event_name = 'dungeon_clear' THEN e.duration_ms / 1000.0 END) AS avg_clear_time,
-- 重复挑战次数
AVG(repeat_attempts) AS avg_attempts
FROM dim_dungeon d
LEFT JOIN (
SELECT
item_id AS dungeon_id,
user_id,
event_name,
duration_ms,
COUNT(*) OVER (PARTITION BY user_id, item_id ORDER BY event_time) AS repeat_attempts
FROM dwd_event_fact
WHERE event_name IN ('dungeon_entry', 'dungeon_clear')
AND dt >= '20260701'
) e ON d.dungeon_id = e.dungeon_id
GROUP BY d.dungeon_id, d.dungeon_name, d.difficulty, d.recommended_level
ORDER BY pass_rate ASC
LIMIT 20;社交分析
-- 社交行为分析
SELECT
dt,
-- 好友相关
AVG(friend_count) AS avg_friend_count,
-- 公会活跃
COUNT(DISTINCT CASE WHEN guild_activity > 0 THEN user_id END) AS guild_active_users,
AVG(guild_activity) AS avg_guild_activity,
-- 组队
COUNT(DISTINCT CASE WHEN party_battle_count > 0 THEN user_id END) AS party_users,
SUM(party_battle_count) AS total_party_battles,
-- 分享传播
COUNT(DISTINCT CASE WHEN share_count > 0 THEN user_id END) AS sharing_users,
SUM(share_count) AS total_shares
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt
ORDER BY dt;4.2 经济系统分析
4.2.1 货币流通分析
-- 货币产出与消耗分析
SELECT
dt,
'金币' AS currency_type,
SUM(gold_income) AS total_output,
SUM(gold_expense) AS total_consumption,
SUM(gold_income) - SUM(gold_expense) AS net_supply,
CASE WHEN SUM(gold_expense) > 0
THEN ROUND(SUM(gold_income) / SUM(gold_expense), 4)
ELSE NULL END AS output_consumption_ratio
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt
UNION ALL
SELECT
dt,
'钻石' AS currency_type,
SUM(diamond_income) AS total_output,
SUM(diamond_expense) AS total_consumption,
SUM(diamond_income) - SUM(diamond_expense) AS net_supply,
CASE WHEN SUM(diamond_expense) > 0
THEN ROUND(SUM(diamond_income) / SUM(diamond_expense), 4)
ELSE NULL END AS output_consumption_ratio
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt
ORDER BY dt, currency_type;
-- 通胀监控(货币总量变化)
SELECT
dt,
SUM(gold_balance) AS total_gold_supply,
SUM(diamond_balance) AS total_diamond_supply,
ROUND((SUM(gold_balance) - LAG(SUM(gold_balance)) OVER(ORDER BY dt))
/ LAG(SUM(gold_balance)) OVER(ORDER BY dt) * 100, 4) AS gold_inflation_rate,
ROUND((SUM(diamond_balance) - LAG(SUM(diamond_balance)) OVER(ORDER BY dt))
/ LAG(SUM(diamond_balance)) OVER(ORDER BY dt) * 100, 4) AS diamond_inflation_rate
FROM dws_player_daily_agg
WHERE dt >= '20260701'
GROUP BY dt;4.2.2 交易市场与拍卖行
-- 玩家交易分析
SELECT
dt,
-- 交易活跃度
COUNT(DISTINCT seller_id) AS seller_count,
COUNT(DISTINCT buyer_id) AS buyer_count,
COUNT(DISTINCT trade_id) AS trade_count,
-- 交易规模
SUM(amount) AS total_trade_amount,
AVG(amount) AS avg_trade_amount,
-- 交易手续费(系统回收货币)
SUM(fee_amount) AS total_fee_collected,
-- 热销道具
item_id,
COUNT(*) AS trade_times
FROM fact_player_trade
WHERE dt >= '20260701'
GROUP BY dt, item_id
ORDER BY dt, trade_times DESC;
-- 交易对价分析(商品价格趋势)
SELECT
item_id,
dt,
AVG(unit_price) AS avg_price,
PERCENTILE(unit_price, 0.5) AS median_price,
MIN(unit_price) AS min_price,
MAX(unit_price) AS max_price,
COUNT(*) AS trade_count,
STDDEV(unit_price) AS price_volatility
FROM fact_player_trade
WHERE dt >= '20260701'
AND item_id = 'target_item'
GROUP BY item_id, dt
ORDER BY dt;4.2.3 经济健康度
class EconomicHealthAnalyzer:
"""经济健康度分析"""
def analyze_economy(self, agg_data):
"""
分析经济系统健康度
核心指标:
- GDP (总产出): 所有货币产出之和
- NNP (净产出): GDP - 折旧(系统回收)
- M2 (货币总量): 流通中所有货币
- 货币流通速度: GDP / M2
"""
total_output = agg_data['gold_income'].sum() + agg_data['diamond_income'].sum()
total_consumption = agg_data['gold_expense'].sum() + agg_data['diamond_expense'].sum()
total_supply = agg_data['gold_balance'].sum() + agg_data['diamond_balance'].sum()
gdp = total_output
nnp = total_output - agg_data['system_fee'].sum()
money_velocity = gdp / max(total_supply, 1)
# 通胀系数
supply_growth_rate = agg_data['total_supply'].pct_change().mean()
# 产出消耗比(1.0 为均衡)
output_consumption_ratio = total_output / max(total_consumption, 1)
# Gini 系数(财富集中度)
gini = self.calculate_gini(agg_data['gold_balance'].values)
return {
'gdp': float(gdp),
'nnp': float(nnp),
'm2_money_supply': float(total_supply),
'money_velocity': float(money_velocity),
'inflation_rate': float(supply_growth_rate),
'output_consumption_ratio': float(output_consumption_ratio),
'gini_coefficient': float(gini),
'health_score': self.calculate_health_score(
output_consumption_ratio, money_velocity, gini
)
}
def calculate_gini(self, values):
"""计算基尼系数"""
sorted_values = np.sort(values)
n = len(sorted_values)
cum_values = np.cumsum(sorted_values)
return (n + 1 - 2 * np.sum(cum_values) / cum_values[-1]) / n
def calculate_health_score(self, oc_ratio, velocity, gini):
"""计算经济健康度评分 (0-100)"""
score = 100
# 产出消耗比惩罚
if oc_ratio > 2.0:
score -= 20 # 产出过多,通胀风险
elif oc_ratio < 0.5:
score -= 15 # 消耗过多,通缩风险
# 货币流通速度惩罚
if velocity < 0.1:
score -= 15 # 货币沉淀严重
elif velocity > 5:
score -= 10 # 货币流转过快
# 基尼系数惩罚(>0.5 表示财富高度集中)
if gini > 0.6:
score -= 25
elif gini > 0.5:
score -= 15
elif gini > 0.4:
score -= 5
return max(0, score)4.3 反作弊数据分析
4.3.1 异常检测算法
# Isolation Forest 孤立森林算法
import numpy as np
from sklearn.ensemble import IsolationForest
def isolation_forest_anomaly_detection(player_features):
"""
使用孤立森林检测异常玩家
特征维度:
- daily_online_time: 日在线时长
- action_frequency: 操作频率
- level_growth_rate: 等级增长速度
- resource_income_rate: 资源获取速度
- gold_growth: 金币增长速度
- diamond_growth: 钻石增长速度
- dungeon_pass_rate: 关卡通过率
- pvp_win_rate: PVP胜率
"""
model = IsolationForest(
n_estimators=200,
max_samples='auto',
contamination=0.01, # 预期异常比例
random_state=42,
n_jobs=-1
)
# 训练模型
model.fit(player_features)
# 预测异常分数(越低越异常)
scores = model.score_samples(player_features)
predictions = model.predict(player_features) # -1: 异常, 1: 正常
# 返回异常玩家
anomalies = []
for i, (pred, score) in enumerate(zip(predictions, scores)):
if pred == -1:
anomalies.append({
'player_index': i,
'anomaly_score': float(score),
'feature_contribution': get_feature_importance(
model, player_features[i:i+1]
)
})
return anomalies
def get_feature_importance(model, sample):
"""计算各特征对异常判定的贡献"""
feature_names = [
'online_time', 'action_freq', 'level_growth', 'resource_rate',
'gold_growth', 'diamond_growth', 'dungeon_pass_rate', 'pvp_win_rate'
]
# 使用路径长度计算特征贡献
contributions = {}
for i, name in enumerate(feature_names):
contributions[name] = float(sample[0, i])
return contributionsLOF (局部离群因子) / DBSCAN 聚类异常检测
from sklearn.neighbors import LocalOutlierFactor
from sklearn.cluster import DBSCAN
from sklearn.preprocessing import StandardScaler
def lof_anomaly_detection(player_features):
"""LOF 局部离群因子检测"""
scaler = StandardScaler()
scaled_features = scaler.fit_transform(player_features)
lof = LocalOutlierFactor(
n_neighbors=20,
contamination=0.05,
novelty=False
)
predictions = lof.fit_predict(scaled_features)
lof_scores = -lof.negative_outlier_factor_ # 越大越异常
return [i for i, p in enumerate(predictions) if p == -1], lof_scores
def dbscan_cluster_anomaly(player_features, eps=0.5, min_samples=5):
"""DBSCAN 聚类异常检测"""
scaler = StandardScaler()
scaled_features = scaler.fit_transform(player_features)
dbscan = DBSCAN(eps=eps, min_samples=min_samples, n_jobs=-1)
clusters = dbscan.fit_predict(scaled_features)
# -1 表示噪声点(异常点)
anomalies = np.where(clusters == -1)[0]
# 统计各聚类大小
cluster_counts = pd.Series(clusters).value_counts()
print(f"聚类分布:\n{cluster_counts}")
# 小型聚类也可能是异常(团伙)
small_clusters = cluster_counts[cluster_counts < min_samples].index.tolist()
small_cluster_anomalies = np.where(np.isin(clusters, small_clusters))[0]
return {
'noise_anomalies': anomalies.tolist(),
'small_cluster_anomalies': small_cluster_anomalies.tolist(),
'cluster_labels': clusters.tolist()
}4.3.2 设备与 IP 聚集分析
-- 同设备多账号检测
SELECT
device_id,
COUNT(DISTINCT user_id) AS account_count,
COUNT(DISTINCT role_id) AS role_count,
COUNT(DISTINCT server_id) AS server_count,
MIN(register_date) AS first_register,
MAX(register_date) AS last_register,
COLLECT_LIST(user_id) AS user_list
FROM dim_player
WHERE device_id IS NOT NULL
AND device_id != ''
GROUP BY device_id
HAVING COUNT(DISTINCT user_id) > 5
ORDER BY account_count DESC;
-- 同 IP 多账号检测
SELECT
register_ip,
COUNT(DISTINCT user_id) AS account_count,
COUNT(DISTINCT device_id) AS device_count,
MIN(register_time) AS first_register,
MAX(register_time) AS last_register,
-- 计算注册间隔(毫秒)
MAX(register_time) - MIN(register_time) AS register_span_ms
FROM dim_player
WHERE register_ip IS NOT NULL
AND register_ip NOT IN ('127.0.0.1', 'localhost')
GROUP BY register_ip
HAVING COUNT(DISTINCT user_id) > 10
ORDER BY account_count DESC;
-- 设备与 IP 交叉聚集(发现团伙)
SELECT
device_id,
register_ip,
COUNT(DISTINCT user_id) AS account_count,
ARRAY_JOIN(COLLECT_SET(register_channel), ',') AS channels
FROM dim_player
WHERE dt = '${dt}'
GROUP BY device_id, register_ip
HAVING COUNT(DISTINCT user_id) > 3;4.3.3 行为序列异常
def detect_behavior_sequence_anomaly(behavior_logs):
"""
检测行为序列异常
正常玩家行为模式示例:
login -> menu -> dungeon_select -> dungeon_start -> dungeon_play -> dungeon_end
login -> mail_check -> shop -> purchase -> inventory
工作室/脚本行为模式示例(过于规律):
login -> dungeon_start -> dungeon_end (循环,间隔几乎相同)
"""
# 提取行为序列
sequences = defaultdict(list)
for log in behavior_logs:
sequences[log.user_id].append({
'event': log.event_name,
'time': log.event_time,
'interval': None # 和上一个事件的间隔
})
anomalies = []
for user_id, seq in sequences.items():
if len(seq) < 10:
continue
# 计算事件间隔
for i in range(1, len(seq)):
seq[i]['interval'] = seq[i]['time'] - seq[i-1]['time']
intervals = [s['interval'] for s in seq[1:] if s['interval'] is not None]
if not intervals:
continue
# 检测指标
interval_std = np.std(intervals)
interval_mean = np.mean(intervals)
cv = interval_std / max(interval_mean, 1) # 变异系数
# 序列熵(衡量行为多样性)
event_types = [s['event'] for s in seq]
entropy = calculate_entropy(event_types)
# 循环检测(相同子序列重复出现)
cycle_score = detect_cycles(event_types)
# 综合异常评分
anomaly_score = 0
if cv < 0.1 and len(intervals) > 20:
anomaly_score += 40 # 间隔过于均匀
if entropy < 1.5:
anomaly_score += 30 # 行为过于单一
if cycle_score > 0.8:
anomaly_score += 30 # 存在明显循环模式
if anomaly_score > 50:
anomalies.append({
'user_id': user_id,
'anomaly_score': anomaly_score,
'sequence_length': len(seq),
'cv': float(cv),
'entropy': float(entropy),
'cycle_score': float(cycle_score)
})
return sorted(anomalies, key=lambda x: x['anomaly_score'], reverse=True)
def calculate_entropy(events):
"""计算序列的信息熵"""
from collections import Counter
import math
counter = Counter(events)
total = len(events)
entropy = 0
for count in counter.values():
p = count / total
entropy -= p * math.log2(p)
return entropy
def detect_cycles(events, min_cycle_length=3, max_cycle_length=10):
"""检测行为序列中是否存在循环模式"""
n = len(events)
max_score = 0
for cycle_len in range(min_cycle_length, min(max_cycle_length, n // 2)):
match_count = 0
total_pairs = 0
for i in range(0, n - cycle_len * 2, cycle_len):
seq1 = events[i:i+cycle_len]
seq2 = events[i+cycle_len:i+cycle_len*2]
if len(seq1) == len(seq2):
total_pairs += 1
if seq1 == seq2:
match_count += 1
if total_parts > 0:
score = match_count / total_pairs
max_score = max(max_score, score)
return max_score4.3.4 图分析团伙检测
import networkx as nx
from community import community_louvain
class CollusionDetection:
"""图分析团伙检测"""
def build_player_graph(self, player_relations):
"""
构建玩家关系图
边权重基于:
- 同IP登录
- 同设备登录
- 频繁交易
- 同公会
- 同注册渠道
"""
G = nx.Graph()
for rel in player_relations:
player_a = rel['player_a']
player_b = rel['player_b']
weight = rel.get('weight', 1.0)
relation_type = rel.get('type', 'unknown')
G.add_edge(player_a, player_b, weight=weight, type=relation_type)
return G
def detect_communities(self, G):
"""
使用 Louvain 算法进行社区发现
同一个社区内的玩家可能存在团伙关系
"""
# Louvain 社区发现
partition = community_louvain.best_partition(G, weight='weight')
# 统计社区大小
communities = {}
for node, community_id in partition.items():
if community_id not in communities:
communities[community_id] = []
communities[community_id].append(node)
# 过滤小社区
suspicious_communities = {
cid: members
for cid, members in communities.items()
if len(members) >= 3 # 3人以上可能构成团伙
}
return suspicious_communities, partition
def find_connected_components(self, G, min_size=3):
"""
连通子图分析
可以找到通过共享设备/IP等连接的玩家团伙
"""
# 使用 IP 和设备搭建二分图
components = list(nx.connected_components(G))
# 过滤规模过小的组件
large_components = [c for c in components if len(c) >= min_size]
return large_components
def identify_ring_leader(self, G, partition):
"""
识别团伙核心(庄家账号)
通过中心性分析找到团伙中的关键节点
"""
leaders = []
for community_id, members in partition.items():
if len(members) < 3:
continue
# 计算子图
subgraph = G.subgraph(members)
# 度中心性:谁连接的团伙成员最多
degree_centrality = nx.degree_centrality(subgraph)
# 介数中心性:谁是团伙内的交易枢纽
betweenness_centrality = nx.betweenness_centrality(subgraph, weight='weight')
# 综合评分
for member in members:
score = (
degree_centrality.get(member, 0) * 0.4 +
betweenness_centrality.get(member, 0) * 0.4 +
G.degree(member, weight='weight') / max(G.degree(), 1) * 0.2
)
leaders.append({
'player_id': member,
'community_id': community_id,
'centrality_score': score,
'degree_centrality': degree_centrality.get(member, 0),
'betweenness_centrality': betweenness_centrality.get(member, 0)
})
# 按综合评分排序
leaders.sort(key=lambda x: x['centrality_score'], reverse=True)
return leaders[:20] # 返回前20个可疑庄家
def detect_offline_trading(self, trade_graph):
"""
检测线下交易(RMT)
特征:
- 两个账号之间频繁且单向的大额交易
- 交易双方 IP 地址不同
- 交易时间集中在特定时段
- 涉及大量金币/高价值道具
"""
suspicious_trades = []
for edge in trade_graph.edges(data=True):
player_a, player_b, data = edge
trade_count = data.get('trade_count', 0)
total_amount = data.get('total_amount', 0)
direction_ratio = data.get('direction_ratio', 0.5)
# 单向大额交易
if abs(direction_ratio - 1.0) < 0.1 and total_amount > 10000:
suspicious_trades.append({
'player_a': player_a,
'player_b': player_b,
'trade_count': trade_count,
'total_amount': total_amount,
'direction_ratio': direction_ratio,
'suspicion_reason': '单向大额交易'
})
return suspicious_trades4.3.5 黑产识别与线下交易检测
-- 线下交易(RMT)检测 SQL
WITH trade_stats AS (
SELECT
seller_id,
buyer_id,
COUNT(*) AS trade_count,
SUM(amount) AS total_amount,
-- 交易方向性:卖方收入占比
SUM(CASE WHEN seller_id = 'target' THEN amount ELSE 0 END) AS seller_income,
SUM(CASE WHEN buyer_id = 'target' THEN amount ELSE 0 END) AS buyer_expense,
-- 交易IP
COUNT(DISTINCT seller_ip) AS seller_ip_count,
COUNT(DISTINCT buyer_ip) AS buyer_ip_count,
-- 交易时间集中度(夜间交易占比)
SUM(CASE WHEN HOUR(trade_time) BETWEEN 0 AND 6 THEN 1 ELSE 0 END) AS night_trades
FROM fact_player_trade
WHERE dt >= '20260701'
GROUP BY seller_id, buyer_id
HAVING COUNT(*) >= 10 -- 频繁交易
AND SUM(amount) > 10000 -- 大额
)
SELECT
seller_id,
buyer_id,
trade_count,
total_amount,
-- 交易方向失衡系数(接近1或0表示单向)
ROUND(ABS(seller_income - buyer_expense) / (seller_income + buyer_expense), 4) AS direction_imbalance,
-- 夜间交易比例
ROUND(night_trades * 1.0 / trade_count, 4) AS night_trade_ratio,
-- 判定
CASE
WHEN direction_imbalance > 0.8 AND total_amount > 50000 THEN '严重可疑'
WHEN direction_imbalance > 0.6 AND night_trade_ratio > 0.3 THEN '高度可疑'
WHEN total_amount > 100000 THEN '大额可疑'
ELSE '正常'
END AS suspicion_level
FROM trade_stats
ORDER BY total_amount DESC;