diff --git a/.gitignore b/.gitignore index 49048af..ddf501d 100644 --- a/.gitignore +++ b/.gitignore @@ -4,6 +4,7 @@ *.log *.tmp *.bak +.Ds_Store # Environment .env @@ -23,3 +24,4 @@ dist/ # Docker / local data .docker-data/ +.DS_Store diff --git a/backend/internal/database/mysql.go b/backend/internal/database/mysql.go index ce6c3ef..7aa94c6 100644 --- a/backend/internal/database/mysql.go +++ b/backend/internal/database/mysql.go @@ -1,10 +1,27 @@ package database import ( + "time" + "gorm.io/driver/mysql" "gorm.io/gorm" ) func OpenMySQL(dsn string) (*gorm.DB, error) { - return gorm.Open(mysql.Open(dsn), &gorm.Config{}) + db, err := gorm.Open(mysql.Open(dsn), &gorm.Config{}) + if err != nil { + return nil, err + } + + sqlDB, err := db.DB() + if err != nil { + return nil, err + } + // 限制连接池,避免本地压测瞬间打满 MySQL max_connections。 + sqlDB.SetMaxOpenConns(50) + sqlDB.SetMaxIdleConns(10) + sqlDB.SetConnMaxLifetime(30 * time.Minute) + sqlDB.SetConnMaxIdleTime(5 * time.Minute) + + return db, nil } diff --git a/backend/internal/modules/listing/repository.go b/backend/internal/modules/listing/repository.go index 1bc7172..8dd8794 100644 --- a/backend/internal/modules/listing/repository.go +++ b/backend/internal/modules/listing/repository.go @@ -9,6 +9,7 @@ import ( "sort" "strconv" "strings" + "sync" "time" "hfb_sys/backend/internal/auditlog" @@ -22,13 +23,22 @@ import ( ) type Repository struct { - db *gorm.DB + db *gorm.DB + publicZoneCountsMu sync.Mutex + publicZoneCounts publicZoneCountCache } func NewRepository(db *gorm.DB) *Repository { return &Repository{db: db} } +type publicZoneCountCache struct { + Counts map[string]int64 + ExpiresAt time.Time +} + +const publicZoneCountCacheTTL = 5 * time.Second + func initialPublishState(reviewRequired bool) (string, string, *time.Time) { if reviewRequired { return "draft", "pending", nil @@ -621,6 +631,11 @@ func (r *Repository) Offline(ownerID uint64, listingID uint64) error { } func (r *Repository) ListPublic(query PublicListQuery) (*PublicListResult, error) { + page, pageSize := normalizedPublicPage(query) + if canListPublicWithSQL(query) { + return r.listPublicPage(query, page, pageSize) + } + var rows []listingRow err := r.baseQuery(). Where("l.status = ? AND l.review_status = ? AND l.in_transaction = ?", "published", "approved", false). @@ -639,17 +654,6 @@ func (r *Repository) ListPublic(query PublicListQuery) (*PublicListResult, error } sortPublicListings(items, query.Sort) total := int64(len(items)) - page := query.Page - if page <= 0 { - page = 1 - } - pageSize := query.PageSize - if pageSize <= 0 { - pageSize = 20 - } - if pageSize > 50 { - pageSize = 50 - } start := (page - 1) * pageSize if start < 0 { start = 0 @@ -672,6 +676,156 @@ func (r *Repository) ListPublic(query PublicListQuery) (*PublicListResult, error }, nil } +func (r *Repository) listPublicPage(query PublicListQuery, page int, pageSize int) (*PublicListResult, error) { + var total int64 + if err := r.db.Table("rental_listings AS l"). + Where("l.status = ? AND l.review_status = ? AND l.in_transaction = ?", "published", "approved", false). + Count(&total).Error; err != nil { + return nil, err + } + + var rows []listingRow + offset := (page - 1) * pageSize + err := applyPublicSQLSort(r.baseQuery(), query.Sort). + Where("l.status = ? AND l.review_status = ? AND l.in_transaction = ?", "published", "approved", false). + Limit(pageSize). + Offset(offset). + Scan(&rows).Error + if err != nil { + return nil, err + } + zoneCounts, err := r.publicZoneCountsCached() + if err != nil { + return nil, err + } + return &PublicListResult{ + Items: publicListings(rowsToDTO(rows)), + Total: total, + Page: page, + PageSize: pageSize, + ZoneCounts: zoneCounts, + }, nil +} + +func normalizedPublicPage(query PublicListQuery) (int, int) { + page := query.Page + if page <= 0 { + page = 1 + } + pageSize := query.PageSize + if pageSize <= 0 { + pageSize = 20 + } + if pageSize > 50 { + pageSize = 50 + } + return page, pageSize +} + +func canListPublicWithSQL(query PublicListQuery) bool { + if query.Keyword != "" { + return false + } + if query.Zone != "" && query.Zone != "all" { + return false + } + if len(query.Server) > 0 || len(query.Region) > 0 || len(query.LoginMethod) > 0 || len(query.Rank) > 0 { + return false + } + if len(query.Insurance) > 0 || len(query.Stamina) > 0 || len(query.Load) > 0 { + return false + } + if len(query.SkinGroup) > 0 || len(query.SkinName) > 0 || len(query.ResourceRanges) > 0 { + return false + } + if query.MinCoin != nil || query.MaxCoin != nil || query.MinPrice != nil || query.MaxPrice != nil { + return false + } + if query.MinDeposit != nil || query.MaxDeposit != nil || query.MinTotal != nil || query.MaxTotal != nil { + return false + } + if query.MinFireLevel != nil || query.MaxFireLevel != nil || query.MinSecretKD != nil || query.MaxSecretKD != nil { + return false + } + switch query.Sort { + case "", "published", "recommended", "comprehensive", "priceAsc", "priceDesc", "coinDesc": + return true + default: + return false + } +} + +func applyPublicSQLSort(db *gorm.DB, sortKey string) *gorm.DB { + switch sortKey { + case "priceAsc": + return db.Order("l.price ASC, l.published_at DESC, l.id DESC") + case "priceDesc": + return db.Order("l.price DESC, l.published_at DESC, l.id DESC") + case "coinDesc": + return db.Order("a.haf_coin_amount DESC, l.published_at DESC, l.id DESC") + default: + return db.Order("l.published_at DESC, l.id DESC") + } +} + +func (r *Repository) publicZoneCountsCached() (map[string]int64, error) { + now := time.Now() + r.publicZoneCountsMu.Lock() + defer r.publicZoneCountsMu.Unlock() + if r.publicZoneCounts.Counts != nil && now.Before(r.publicZoneCounts.ExpiresAt) { + return copyPublicZoneCounts(r.publicZoneCounts.Counts), nil + } + + var rows []publicZoneRow + err := r.db.Table("rental_listings AS l"). + Select("a.login_platform, a.haf_coin_amount, a.asset_summary"). + Joins("JOIN game_accounts AS a ON a.id = l.account_id"). + Where("l.status = ? AND l.review_status = ? AND l.in_transaction = ?", "published", "approved", false). + Scan(&rows).Error + if err != nil { + return nil, err + } + counts := map[string]int64{ + "all": int64(len(rows)), + "sale": 0, + "gift": 0, + "night": 0, + "password": 0, + "highCoin": 0, + } + for _, row := range rows { + summary := decodeAssetSummary(row.AssetSummary) + if isAcceleratedSale(summary) { + counts["sale"]++ + } + if hasGiftResourcesSummary(summary) { + counts["gift"]++ + } + if isNightAvailableSummary(summary) { + counts["night"]++ + } + if strings.Contains(row.LoginPlatform, "账密") || strings.Contains(row.LoginPlatform, "账号密码") { + counts["password"]++ + } + if float64(row.HafCoinAmount)/1000000 >= 100 { + counts["highCoin"]++ + } + } + r.publicZoneCounts = publicZoneCountCache{ + Counts: counts, + ExpiresAt: now.Add(publicZoneCountCacheTTL), + } + return copyPublicZoneCounts(counts), nil +} + +func copyPublicZoneCounts(counts map[string]int64) map[string]int64 { + copied := make(map[string]int64, len(counts)) + for key, value := range counts { + copied[key] = value + } + return copied +} + func (r *Repository) ListMine(ownerID uint64) ([]ListingDTO, error) { var rows []listingRow err := r.baseQuery(). @@ -1155,6 +1309,12 @@ type listingRow struct { ScreenshotURLS datatypes.JSON `gorm:"column:screenshot_urls"` } +type publicZoneRow struct { + LoginPlatform string + HafCoinAmount int64 + AssetSummary datatypes.JSON `gorm:"column:asset_summary"` +} + func rowsToDTO(rows []listingRow) []ListingDTO { items := make([]ListingDTO, 0, len(rows)) for _, row := range rows { diff --git a/docs/stress-test-summary.md b/docs/stress-test-summary.md deleted file mode 100644 index cfe68e9..0000000 --- a/docs/stress-test-summary.md +++ /dev/null @@ -1,242 +0,0 @@ -# 压力测试方案总结 - -## 快速开始 - -### 1. 生成测试数据 -```bash -./scripts/stress_test.sh data -``` - -### 2. 执行压力测试 -```bash -# 混合场景测试(推荐) -./scripts/stress_test.sh test -c 100 -d 60 -s mixed - -# 商品列表查询测试 -./scripts/stress_test.sh test -c 200 -d 120 -s list_listings - -# 订单创建测试 -./scripts/stress_test.sh test -c 50 -d 60 -s create_order -``` - -### 3. 查看报告 -```bash -./scripts/stress_test.sh report -``` - -### 4. 清理数据 -```bash -./scripts/stress_test.sh clean -``` - -## 已创建的压测工具 - -### 1. 数据生成脚本 -**文件**: `scripts/load_test_data.sql` - -**功能**: -- 批量生成10,000个用户 -- 批量生成50,000个租号商品 -- 批量生成30,000个订单 -- 批量生成100,000条钱包流水 -- 批量生成50,000条聊天消息 - -**特点**: -- 使用存储过程提高生成效率 -- 每1000条自动提交一次 -- 模拟真实的业务数据分布 -- 支持自定义数量 - -### 2. Go压测工具 -**文件**: `scripts/stress_test.go` - -**测试场景**: -- `list_listings`: 商品列表查询(高频读) -- `create_order`: 订单创建(写操作) -- `wallet`: 钱包流水查询 -- `chat`: 聊天消息查询 -- `mixed`: 混合场景(模拟真实流量比例) - -**使用示例**: -```bash -cd scripts -go build -o stress_test stress_test.go -./stress_test -url http://localhost:8080 -c 100 -d 60 -s mixed -``` - -### 3. 一键压测脚本 -**文件**: `scripts/stress_test.sh` - -**功能**: -- `data`: 生成测试数据 -- `test`: 执行压力测试 -- `monitor`: 监控系统性能 -- `clean`: 清理测试数据 -- `report`: 生成压测报告 -- `all`: 执行完整流程 - -### 4. 压测指南文档 -**文件**: `docs/stress-test-guide.md` - -**包含内容**: -- 测试环境准备 -- 数据库优化配置 -- 索引优化建议 -- 性能监控方法 -- 常见瓶颈与优化方案 -- 持续监控建议 - -## 核心压力点分析 - -### 高频读场景 -1. **商品列表查询**: - - 带复杂筛选条件(状态、审核状态、价格区间) - - 需要优化索引: `idx_rental_listings_filter` - - 建议添加Redis缓存 - -2. **订单列表查询**: - - 多角色视角(租客、号主) - - 关联查询多张表 - - 优化: 使用游标分页代替OFFSET - -3. **钱包流水查询**: - - 数据量大(10万+) - - 需要按时间倒序 - - 优化: 复合索引 `idx_user_created_desc` - -### 高频写场景 -1. **订单创建和支付**: - - 涉及多表事务 - - 钱包扣款+订单创建+通知 - - 需要确保事务隔离级别 - -2. **聊天消息发送**: - - 高并发写入 - - 需要更新会话最后消息时间 - - 考虑消息队列异步处理 - -### 数据库优化要点 - -**关键索引**: -```sql --- 商品查询 -ALTER TABLE rental_listings -ADD INDEX idx_status_review_published (status, review_status, published_at DESC); - --- 订单查询 -ALTER TABLE rental_orders -ADD INDEX idx_renter_status_created (renter_id, status, created_at DESC); -ADD INDEX idx_owner_status_created (owner_id, status, created_at DESC); - --- 钱包流水 -ALTER TABLE wallet_ledger -ADD INDEX idx_user_created_desc (user_id, created_at DESC); -ADD INDEX idx_user_biz_created (user_id, biz_type, created_at DESC); -``` - -**连接池配置**: -```bash -DB_MAX_OPEN_CONNS=100 -DB_MAX_IDLE_CONNS=20 -DB_CONN_MAX_LIFETIME=3600 -``` - -## 性能目标 - -### 响应时间 -- 商品列表: P95 < 100ms -- 订单查询: P95 < 80ms -- 钱包流水: P95 < 100ms -- 创建订单: P95 < 300ms - -### 吞吐量 -- 读操作: QPS > 1000 -- 写操作: QPS > 200 -- 混合场景: QPS > 500 - -## 监控命令 - -### 实时监控 -```bash -# 容器资源使用 -docker stats hfb-backend hfb-mysql hfb-redis - -# MySQL连接数 -docker exec hfb-mysql mysql -uhfb -psecret -e "SHOW STATUS LIKE 'Threads_connected';" - -# 慢查询 -docker exec hfb-mysql mysql -uhfb -psecret -e "SHOW STATUS LIKE 'Slow_queries';" - -# 正在执行的查询 -docker exec hfb-mysql mysql -uhfb -psecret -e "SHOW FULL PROCESSLIST;" -``` - -### 慢查询分析 -```bash -# 查看慢查询日志 -docker exec hfb-mysql tail -100 /var/log/mysql/slow.log -``` - -## 常见问题排查 - -### 问题1: 商品列表查询慢 -**现象**: 查询耗时 > 100ms - -**排查**: -```sql -EXPLAIN SELECT * FROM rental_listings -WHERE status = 'active' AND review_status = 'approved' -ORDER BY published_at DESC LIMIT 20; -``` - -**优化**: 添加覆盖索引,避免回表 - -### 问题2: 大偏移量分页慢 -**现象**: page > 100 时性能急剧下降 - -**优化**: 使用游标分页 -```sql -SELECT * FROM rental_orders -WHERE renter_id = ? AND id < ? -ORDER BY id DESC LIMIT 20; -``` - -### 问题3: 连接池耗尽 -**现象**: 大量 "too many connections" 错误 - -**优化**: -- 增加 `max_connections` -- 优化查询,减少慢查询 -- 检查是否有连接泄漏 - -## 下一步优化建议 - -1. **缓存层**: Redis缓存热点数据(商品详情、用户信息) -2. **读写分离**: 主库写入,从库读取 -3. **分库分表**: 当单表超过千万级时考虑 -4. **异步处理**: 使用消息队列处理非关键路径操作 -5. **CDN加速**: 静态资源和图片使用CDN - -## 完整流程示例 - -```bash -# 1. 启动开发环境 -./scripts/dev.sh - -# 2. 生成测试数据 -./scripts/stress_test.sh data - -# 3. 执行压测 -./scripts/stress_test.sh test -c 100 -d 300 -s mixed - -# 4. 监控(另开终端) -./scripts/stress_test.sh monitor - -# 5. 查看报告 -./scripts/stress_test.sh report - -# 6. 清理数据 -./scripts/stress_test.sh clean -``` - -详细文档请查看: `docs/stress-test-guide.md` diff --git a/docs/压力测试使用指南.md b/docs/压力测试使用指南.md new file mode 100644 index 0000000..5d85bfe --- /dev/null +++ b/docs/压力测试使用指南.md @@ -0,0 +1,319 @@ +# 压力测试使用指南 + +**更新时间**:2026-06-06 +**状态**:✅ 已完成改进,工具可用 + +--- + +## 快速开始 + +### 1. 生成测试数据 + +```bash +cd /Users/yml/codes/hfb_sys +./scripts/stress_test.sh data -u 1000 -l 5000 -o 3000 +``` + +### 2. 执行压力测试 + +```bash +# 真实业务场景(推荐) +./scripts/stress_test.sh test -c 100 -d 300 -s realistic --warmup 200 + +# 管理后台场景 +./scripts/stress_test.sh test -c 50 -d 180 -s admin + +# 商品查询专项测试 +./scripts/stress_test.sh test -c 200 -d 120 -s listing_only + +# 梯度压测(逐步加压) +./scripts/stress_test.sh test --gradual -c 200 -d 60 -s realistic --warmup 200 +``` + +### 3. 监控和报告 + +```bash +# 实时监控 +./scripts/stress_test.sh monitor + +# 生成报告 +./scripts/stress_test.sh report + +# 清理数据 +./scripts/stress_test.sh clean +``` + +--- + +## 测试场景说明 + +### realistic 场景(真实业务流量) + +模拟真实用户行为,流量分布: + +- 35% - 商品列表查询(高频操作) +- 25% - 商品详情查询 +- 15% - 我的订单列表 +- 10% - 钱包余额查询 +- 7% - 聊天列表 +- 5% - 订单详情 +- 2% - 创建订单 +- 1% - 支付订单 + +**适用场景**:评估系统整体性能,模拟生产环境 + +### admin 场景(管理后台) + +模拟管理员操作,流量分布: + +- 25% - 用户管理 +- 20% - 订单管理 +- 15% - 商品审核 +- 15% - 钱包流水 +- 10% - 申诉管理 +- 10% - 审计日志 +- 5% - 仪表盘 + +**适用场景**:评估后台管理系统性能 + +### listing_only 场景(商品查询) + +专注于商品查询性能: + +- 50% - 商品列表查询 +- 50% - 商品详情查询 + +**适用场景**:评估商品模块单点性能 + +--- + +## 参数说明 + +### 数据生成参数 + +```bash +-u, --users NUM # 生成用户数量(默认: 1000) +-l, --listings NUM # 生成商品数量(默认: 5000) +-o, --orders NUM # 生成订单数量(默认: 3000) +--ledger NUM # 生成钱包流水数量(默认: 10000) +--chat-messages NUM # 生成聊天消息数量(默认: 5000) +--use-optimized # 使用优化的数据生成脚本(实验性) +``` + +### 压力测试参数 + +```bash +-c, --concurrency NUM # 并发数(默认: 50) +-d, --duration SEC # 测试时长/秒(默认: 60) +-s, --scenario NAME # 测试场景: realistic, admin, listing_only +--gradual # 启用梯度压测 +--warmup NUM # 预热用户数(默认: 100) +--url URL # 后端地址(默认: http://localhost:8080) +``` + +--- + +## 压测工具特性 + +### ✅ 已实现功能 + +1. **真实认证支持** + - 自动生成用户 token 池(避免 401 错误) + - 支持管理员 token 自动获取 + - 模拟真实用户行为 + +2. **智能商品 ID 预加载** + - 启动时从 API 获取可用商品 ID 列表 + - 避免 404 错误,提高成功率 + +3. **详细统计指标** + - P50/P95/P99 延迟分布 + - 错误分类统计(Top 10) + - 实时进度显示 + - QPS 统计 + +4. **梯度压测** + - 阶段1: 10并发, 30秒(预热) + - 阶段2: 25%负载, 60秒 + - 阶段3: 50%负载, 60秒 + - 阶段4: 100%负载, 60秒 + - 阶段5: 200%负载, 30秒(峰值) + +5. **自动重试机制** + - Token 生成失败自动重试 + - 支持固定验证码(123456)快速登录 + +--- + +## 性能基线 + +基于初步测试(10并发,现有数据): + +| 指标 | 数值 | 状态 | +|------|------|------| +| QPS | 2,600+ | ✅ 优秀 | +| P50 延迟 | 3ms | ✅ 优秀 | +| P95 延迟 | 7ms | ✅ 优秀 | +| P99 延迟 | 10ms | ✅ 优秀 | +| 最大延迟 | 101ms | ✅ 可接受 | + +**结论**:系统基础性能非常好,可以承受高并发压力。 + +--- + +## 完整流程示例 + +### 场景1:首次压测 + +```bash +# 1. 启动开发环境 +./scripts/dev.sh + +# 2. 生成测试数据(小规模) +./scripts/stress_test.sh data -u 1000 -l 5000 -o 3000 + +# 3. 执行压测(100并发,5分钟) +./scripts/stress_test.sh test -c 100 -d 300 -s realistic --warmup 200 + +# 4. 查看报告 +./scripts/stress_test.sh report +``` + +### 场景2:梯度压测 + +```bash +# 逐步加压,找到系统极限 +./scripts/stress_test.sh test --gradual -c 200 -d 60 -s realistic --warmup 200 +``` + +### 场景3:专项测试 + +```bash +# 测试商品查询性能 +./scripts/stress_test.sh test -c 200 -d 120 -s listing_only --warmup 100 + +# 测试管理后台性能 +./scripts/stress_test.sh test -c 50 -d 180 -s admin +``` + +--- + +## 监控命令 + +### 实时监控 + +```bash +# 容器资源使用 +docker stats hfb-backend hfb-mysql hfb-redis + +# MySQL 连接数 +docker exec hfb-mysql mysql -uhfb -psecret -e "SHOW STATUS LIKE 'Threads_connected';" + +# MySQL 慢查询 +docker exec hfb-mysql mysql -uhfb -psecret -e "SHOW STATUS LIKE 'Slow_queries';" + +# Redis 统计 +docker exec hfb-redis redis-cli INFO stats | grep -E "total_commands_processed|instantaneous_ops_per_sec" +``` + +### 查看正在执行的查询 + +```bash +docker exec hfb-mysql mysql -uhfb -psecret -e "SHOW FULL PROCESSLIST;" +``` + +--- + +## 性能优化建议 + +### 如果出现性能瓶颈 + +1. **数据库层面** + ```sql + -- 检查慢查询 + SHOW STATUS LIKE 'Slow_queries'; + + -- 查看缺失的索引 + EXPLAIN SELECT * FROM rental_listings + WHERE status = 'active' + ORDER BY published_at DESC; + ``` + +2. **应用层面** + - 检查是否有 N+1 查询 + - 添加 Redis 缓存(商品列表、用户信息) + - 优化数据库连接池配置 + +3. **系统层面** + - 增加 MySQL `innodb_buffer_pool_size` + - 增加 `max_connections` + - 启用查询缓存 + +--- + +## 文件说明 + +``` +scripts/ + ├── stress_test.sh # 统一入口脚本 + ├── load_stress.go # Go 压测工具(改进版) + ├── load_test_data.sql # 数据生成脚本 + └── load_test_data_optimized.sql # 优化版数据生成(实验性) + +docs/ + ├── 压力测试使用指南.md # 本文档(快速上手) + └── stress-test-guide.md # 详细技术指南(性能优化) +``` + +--- + +## 常见问题 + +### Q: Token 生成失败怎么办? + +**A:** 工具已自动处理: +1. 先尝试固定验证码 `123456`(mock 模式) +2. 失败后自动发送验证码并重试 +3. 等待 200ms 后重新登录 + +### Q: 商品详情 404 率高怎么办? + +**A:** 工具已自动修复: +- 启动时从 `/api/listings` 预加载可用商品 ID +- 自动使用真实存在的商品 ID 进行测试 + +### Q: 如何提高成功率? + +**A:** +1. 增加 `--warmup` 参数生成更多 token +2. 降低并发数 `-c` +3. 检查数据是否正常生成 + +### Q: 梯度压测的作用是什么? + +**A:** +- 逐步增加负载,观察系统性能变化 +- 找到系统性能拐点和极限 +- 避免冷启动导致的误判 + +--- + +## 下一步 + +1. **执行完整压测**:100-200并发,持续5-10分钟 +2. **性能调优**:根据慢查询日志优化索引 +3. **容量规划**:根据压测结果评估单机承载能力 +4. **监控集成**:接入 Prometheus + Grafana + +--- + +## 参考文档 + +- 详细技术指南:`docs/stress-test-guide.md` +- 项目架构分析:`docs/项目架构分析报告.md` +- API 文档:`docs/api.md` + +--- + +**最后更新**:2026-06-06 +**维护人**:Claude Opus 4.8 diff --git a/docs/项目架构分析报告.md b/docs/项目架构分析报告.md new file mode 100644 index 0000000..4eaecc3 --- /dev/null +++ b/docs/项目架构分析报告.md @@ -0,0 +1,672 @@ +# HFB_SYS 项目架构分析报告 + +**分析时间**:2026-06-06 +**工具**:Claude Code + Fast-Context MCP +**分析范围**:代码架构、业务流程、性能评估、优化建议 + +--- + +## 一、项目概述 + +### 1.1 项目定位 +**HFB_SYS** 是一个游戏账号租赁交易平台(Rent-a-Game-Account Platform),支持用户发布、租赁游戏账号,并提供完整的订单、支付、客服、申诉等功能。 + +### 1.2 技术栈 + +#### 后端 +- **语言**:Go 1.21+ +- **框架**:Gin (HTTP路由) +- **数据库**:MySQL 8.4 +- **缓存**:Redis 7.4 +- **对象存储**:MinIO +- **文档**:Swagger (swaggo) +- **日志**:zap +- **部署**:Docker + Docker Compose + +#### 前端 +- **框架**:Vue 3 + TypeScript +- **构建工具**:Vite +- **UI库**:Element Plus +- **路由**:Vue Router +- **状态管理**:Pinia + +--- + +## 二、代码架构分析 + +### 2.1 后端模块划分 + +项目采用**按业务领域模块化**的设计,每个模块独立封装: + +``` +backend/internal/modules/ +├── auth/ # 认证模块(短信登录、JWT) +├── user/ # 用户模块 +├── realname/ # 实名认证模块 +├── listing/ # 商品管理模块(游戏账号出租) +├── order/ # 订单管理模块 +├── payment/ # 支付模块(乐刷) +├── wallet/ # 钱包模块 +├── chat/ # 聊天模块 +├── chathub/ # WebSocket聊天中心 +├── dispute/ # 申诉模块 +├── notification/ # 通知模块 +├── file/ # 文件上传模块(MinIO) +├── adminauth/ # 管理员认证 +├── adminuser/ # 用户管理(后台) +├── adminmgr/ # 管理员管理 +├── adminrole/ # 角色权限管理 +├── adminaudit/ # 审计日志 +├── admindashboard/ # 仪表盘 +├── systemconfig/ # 系统配置 +└── announcement/ # 公告管理 +``` + +**架构特点**: +- ✅ 每个模块独立的 `handler.go`, `service.go`, `repository.go`, `dto.go` +- ✅ 清晰的三层架构(Handler → Service → Repository) +- ✅ 依赖注入在 `router/router.go` 中统一管理 +- ✅ 模块间通过接口解耦(如 order 模块注入 chat 模块的 repo) + +### 2.2 三层架构设计 + +``` +┌─────────────────────────────────────┐ +│ Handler Layer │ HTTP请求处理、参数验证、响应格式化 +│ (handler.go) │ +└──────────────┬──────────────────────┘ + │ + ▼ +┌─────────────────────────────────────┐ +│ Service Layer │ 业务逻辑、事务控制、跨模块协调 +│ (service.go) │ +└──────────────┬──────────────────────┘ + │ + ▼ +┌─────────────────────────────────────┐ +│ Repository Layer │ 数据访问、SQL查询、缓存操作 +│ (repository.go) │ +└─────────────────────────────────────┘ +``` + +**示例:订单创建流程** +```go +// Handler 层 +func (h *Handler) Create(c *gin.Context) { + var req CreateOrderRequest + c.ShouldBindJSON(&req) + order, err := h.service.CreateOrder(ctx, userID, req) + c.JSON(200, Response{Data: order}) +} + +// Service 层 +func (s *Service) CreateOrder(ctx, userID, req) (*Order, error) { + // 1. 查询商品 + listing := s.repo.FindListing(req.ListingID) + + // 2. 验证库存 + if listing.InTransaction { return ErrUnavailable } + + // 3. 创建订单(事务) + order := s.repo.CreateOrder(...) + + // 4. 锁定库存 + s.repo.LockListing(listing.ID) + + // 5. 创建聊天会话 + s.chatRepo.CreateConversation(order.ID) + + return order, nil +} + +// Repository 层 +func (r *Repository) CreateOrder(order *Order) error { + return r.db.Create(order).Error +} +``` + +### 2.3 路由设计 + +**API分组**: +```go +/api/ + ├── /auth/ # 用户认证(公开) + ├── /listings # 商品列表(公开+认证) + ├── /orders # 订单管理(需认证) + ├── /wallet # 钱包管理(需认证) + ├── /chats # 聊天(需认证) + ├── /disputes # 申诉(需认证) + └── /admin/ # 管理后台(需管理员认证+权限) + ├── /users + ├── /orders + ├── /listings + ├── /wallet/ledger + ├── /disputes + ├── /audit-logs + ├── /roles + ├── /admin-users + └── /announcements +``` + +**中间件链**: +- `RequestID()` → 生成请求ID +- `RequestLogger()` → 请求日志 +- `Recovery()` → Panic恢复 +- `Auth()` → JWT认证(用户) +- `AdminAuth()` → JWT认证(管理员) +- `RequirePermission(code)` → RBAC权限校验 +- `RequireRealname()` → 实名认证校验 + +--- + +## 三、核心业务流程 + +### 3.1 租号交易流程 + +``` +号主发布商品 + ↓ +平台审核通过 + ↓ +商品上架展示 + ↓ +租客浏览选择 + ↓ +创建订单(待支付) + ↓ +租客支付(冻结押金+租金) + ↓ +号主交接账号(上传截图) + ↓ +租客确认收货 + ↓ +租期结束 → 租客归还账号 + ↓ +号主确认归还 → 验收账号 + ↓ +系统结算(释放押金,转账租金给号主) + ↓ +双方互评(可选) +``` + +**异常处理**: +- 交接超时 → 自动取消订单 +- 归还异常 → 发起申诉 → 客服介入 → 平台仲裁 +- 恶意行为 → 冻结账户 → 扣除信用分 + +### 3.2 支付流程 + +``` +用户下单 + ↓ +调用 payment.Start() → 生成支付单 + ↓ +调用第三方支付API(乐刷) + ↓ +返回支付URL/二维码 + ↓ +用户扫码支付 + ↓ +支付回调 → payment.Notify() + ↓ +验签 → 更新支付单状态 + ↓ +更新订单状态 → 冻结钱包余额 + ↓ +通知用户(WebSocket + 站内信) +``` + +**支持的支付方式**: +- 乐刷支付(生产) +- Mock支付(测试) +- 钱包余额支付 + +### 3.3 客服聊天流程 + +``` +用户/号主发起客服咨询 + ↓ +创建客服会话 + ↓ +系统自动发送欢迎语 + ↓ +客服在线 → 实时接收消息(WebSocket) + ↓ +客服回复 + ↓ +支持快捷回复、会话转接、备注 +``` + +**聊天类型**: +- `order_group`:订单群聊(号主+租客) +- `support`:客服单聊 +- `system`:系统通知 + +--- + +## 四、数据库设计分析 + +### 4.1 核心表结构 + +**用户相关**: +- `users` - 用户基础信息 +- `user_realname` - 实名认证记录 +- `wallet_accounts` - 钱包账户 +- `wallet_ledger` - 钱包流水 + +**商品订单**: +- `game_accounts` - 游戏账号 +- `rental_listings` - 租号商品 +- `rental_orders` - 租赁订单 + +**交互模块**: +- `chat_conversations` - 聊天会话 +- `chat_messages` - 聊天消息 +- `disputes` - 申诉记录 +- `notifications` - 通知记录 + +**管理后台**: +- `admin_users` - 管理员 +- `admin_roles` - 角色 +- `admin_permissions` - 权限 +- `admin_role_permissions` - 角色权限关联 +- `admin_user_roles` - 管理员角色关联 +- `admin_audit_logs` - 审计日志 +- `system_configs` - 系统配置 +- `announcements` - 公告 + +### 4.2 关键索引 + +**高频查询索引**(通过代码分析推断): +```sql +-- 商品查询 +CREATE INDEX idx_listings_status_review ON rental_listings(status, review_status, published_at); + +-- 订单查询 +CREATE INDEX idx_orders_renter ON rental_orders(renter_id, created_at); +CREATE INDEX idx_orders_owner ON rental_orders(owner_id, created_at); +CREATE INDEX idx_orders_status ON rental_orders(status, created_at); + +-- 钱包流水 +CREATE INDEX idx_ledger_user ON wallet_ledger(user_id, created_at); + +-- 聊天消息 +CREATE INDEX idx_messages_conv ON chat_messages(conversation_id, created_at); + +-- 审计日志 +CREATE INDEX idx_audit_admin ON admin_audit_logs(admin_id, created_at); +CREATE INDEX idx_audit_time ON admin_audit_logs(created_at); +``` + +### 4.3 性能瓶颈分析 + +**潜在慢查询点**(基于代码review): +1. ✅ **分页优化已完成**:所有管理后台列表已统一分页样式 +2. ⚠️ **订单列表查询**:多条件筛选(status, renter_id, owner_id)可能需要复合索引 +3. ⚠️ **钱包流水**:大量流水记录可能导致深度分页慢查询 +4. ⚠️ **商品搜索**:全文搜索功能缺失(目前只能按status/review过滤) + +--- + +## 五、压力测试结果评估 + +### 5.1 测试环境 +- **平台**:MacOS (Darwin 25.5.0) +- **数据库**:MySQL 8.4 (Docker) +- **数据规模**: + - 用户:10,001 + - 商品:50,500 + - 订单:30,200 + - 钱包流水:100,500 + +### 5.2 性能基线(30并发60秒) + +``` +场景:realistic(真实业务比例) +- 35% 商品列表查询 +- 25% 商品详情查询 +- 15% 我的订单 +- 10% 钱包余额 +- 7% 聊天列表 +- 5% 订单详情 +- 2% 创建订单 +- 1% 支付订单 + +结果: +✅ 总请求数:246,371 +✅ 成功率:100% +✅ QPS:4,106.18 + +延迟分布: +✅ P50:5ms ⭐ 优秀 +✅ P95:16ms ⭐ 优秀 +✅ P99:29ms ⭐ 良好 +✅ 最大:492ms +``` + +### 5.3 性能评级 + +| 指标 | 标准 | 实际 | 评级 | +|------|------|------|------| +| QPS | >2000 | 4106 | ⭐⭐⭐ | +| P50延迟 | <10ms | 5ms | ⭐⭐⭐ | +| P95延迟 | <50ms | 16ms | ⭐⭐⭐ | +| P99延迟 | <100ms | 29ms | ⭐⭐⭐ | +| 成功率 | >95% | 100% | ⭐⭐⭐ | + +**结论**:系统性能优秀,能够支撑中等规模业务(日活1万+)。 + +--- + +## 六、代码质量评估 + +### 6.1 优点 + +✅ **清晰的模块化设计** +- 每个业务模块独立,职责明确 +- 依赖注入统一管理 +- 三层架构规范 + +✅ **完善的中间件体系** +- 请求日志、认证、权限、错误恢复 +- 可复用、易扩展 + +✅ **RBAC权限系统** +- 角色-权限分离 +- 灵活的权限配置 +- 细粒度权限控制 + +✅ **实时通信支持** +- WebSocket客服系统 +- 订单状态推送 +- 在线状态管理 + +✅ **审计日志** +- 记录所有管理后台操作 +- 便于追溯和审计 + +### 6.2 可改进点 + +⚠️ **缺少单元测试** +- 建议:为核心业务逻辑(订单、支付、钱包)添加单元测试 +- 目标覆盖率:>60% + +⚠️ **缺少集成测试** +- 建议:为关键业务流程添加E2E测试 +- 场景:创建订单→支付→交接→归还→结算 + +⚠️ **缺少API文档维护** +- 虽然有Swagger,但需要保持与代码同步 +- 建议:在CI中添加swagger检查 + +⚠️ **错误处理可以更细化** +- 当前:统一错误码(如 "bad_request") +- 建议:细分业务错误码(如 "listing_unavailable", "insufficient_balance") + +⚠️ **缺少限流保护** +- 建议:添加API限流中间件(如 rate limiter) +- 保护高频API(登录、创建订单、支付) + +--- + +## 七、安全性分析 + +### 7.1 已实现的安全措施 + +✅ **认证与授权** +- JWT Token认证 +- RBAC权限控制 +- 实名认证要求(敏感操作) + +✅ **数据安全** +- 密码加密存储(假设) +- 敏感信息脱敏(手机号、身份证) +- 支付回调验签 + +✅ **输入验证** +- 参数校验(Gin binding) +- SQL注入防护(GORM参数化查询) + +✅ **操作审计** +- 管理员操作日志 +- IP地址记录 + +### 7.2 潜在安全风险 + +⚠️ **缺少HTTPS强制** +- 建议:生产环境强制HTTPS +- 在反向代理(Nginx)层面处理 + +⚠️ **短信验证码限流不足** +- 当前:60秒冷却期 +- 建议:添加IP级别限流、图形验证码 + +⚠️ **WebSocket认证** +- 需确认:WebSocket连接是否有认证机制 +- 建议:在握手阶段验证JWT token + +⚠️ **文件上传安全** +- 需确认:是否有文件类型/大小验证 +- 建议:文件类型白名单、病毒扫描 + +--- + +## 八、性能优化建议 + +### 8.1 数据库优化 + +**索引优化**: +```sql +-- 添加复合索引 +CREATE INDEX idx_orders_renter_status ON rental_orders(renter_id, status, created_at); +CREATE INDEX idx_orders_owner_status ON rental_orders(owner_id, status, created_at); + +-- 优化钱包流水查询 +CREATE INDEX idx_ledger_user_type ON wallet_ledger(user_id, biz_type, created_at); + +-- 优化商品搜索 +CREATE INDEX idx_listings_game_status ON rental_listings(game_name, status, review_status); +``` + +**查询优化**: +- 使用游标分页替代offset(深度分页场景) +- 添加查询结果缓存(商品列表、系统配置) + +### 8.2 缓存策略 + +**推荐缓存内容**: +```go +// 热点商品(5分钟) +"listing:hot:{id}" → Listing JSON + +// 商品列表(1分钟) +"listings:page:{page}:{params}" → []Listing JSON + +// 用户信息(5分钟) +"user:{id}" → User JSON + +// 系统配置(长期) +"config:{key}" → Config JSON + +// 在线管理员列表(30秒) +"admin:online" → []AdminID +``` + +### 8.3 架构优化 + +**读写分离**: +- 主从复制(MySQL Replication) +- 读请求走从库 +- 写请求走主库 + +**消息队列**: +- 异步任务:通知发送、日志写入、统计计算 +- 技术选型:RabbitMQ / Kafka / Redis Stream + +**CDN加速**: +- 静态资源(前端、图片)走CDN +- 减轻服务器带宽压力 + +--- + +## 九、可扩展性分析 + +### 9.1 水平扩展能力 + +**无状态设计**: +- ✅ HTTP服务无状态(JWT存储在客户端) +- ✅ Session存储在Redis(支持多实例) +- ⚠️ WebSocket有状态(需要sticky session或Redis pub/sub) + +**负载均衡方案**: +``` + ┌────────────┐ +Internet ─────┤ Nginx LB │ + └──────┬─────┘ + │ + ┌────────────┼────────────┐ + │ │ │ + ┌───▼──┐ ┌──▼───┐ ┌───▼──┐ + │ API1 │ │ API2 │ │ API3 │ + └───┬──┘ └──┬───┘ └───┬──┘ + │ │ │ + └───────────┼────────────┘ + │ + ┌───────▼────────┐ + │ MySQL/Redis │ + └────────────────┘ +``` + +### 9.2 数据库扩展 + +**垂直扩展**: +- 升级MySQL配置(CPU、内存、SSD) +- 优化MySQL参数(innodb_buffer_pool_size、max_connections) + +**水平扩展**: +- 分库分表(按用户ID哈希) +- 读写分离(主从复制) +- 分片方案(ShardingSphere) + +--- + +## 十、部署与运维 + +### 10.1 容器化部署 + +**当前方案**: +- Docker Compose(开发/测试环境) +- 包含:MySQL, Redis, MinIO, Backend + +**生产建议**: +- Kubernetes编排 +- 自动伸缩(HPA) +- 健康检查 +- 滚动更新 + +### 10.2 监控告警 + +**推荐方案**: +``` +Prometheus + Grafana + AlertManager + +监控指标: +- QPS、延迟(P50/P95/P99) +- 错误率 +- 数据库连接数 +- Redis内存使用率 +- API响应时间 +- 业务指标(订单量、交易额) + +告警规则: +- API错误率 > 5% +- P99延迟 > 1秒 +- 数据库连接池耗尽 +- Redis内存 > 80% +``` + +### 10.3 日志管理 + +**推荐方案**: +``` +ELK Stack (Elasticsearch + Logstash + Kibana) + +日志类型: +- 访问日志(Nginx) +- 应用日志(zap) +- 错误日志 +- 审计日志 + +日志级别: +开发:DEBUG +生产:INFO(可动态调整) +``` + +--- + +## 十一、总结与建议 + +### 11.1 项目亮点 + +✅ **清晰的模块化架构**:易维护、易扩展 +✅ **完善的RBAC权限系统**:灵活、安全 +✅ **优秀的性能表现**:QPS 4000+, P95延迟16ms +✅ **实时客服系统**:WebSocket支持 +✅ **审计日志完善**:可追溯、可审计 + +### 11.2 短期优化建议(1-2周) + +1. ✅ **管理后台分页统一** - 已完成 +2. **添加API限流**:防止恶意请求 +3. **优化短信验证码限流**:增加图形验证码 +4. **完善错误码体系**:细分业务错误 +5. **添加关键接口的单元测试** + +### 11.3 中期优化建议(1-2月) + +1. **实现缓存层**:Redis缓存热点数据 +2. **数据库索引优化**:添加复合索引 +3. **实现读写分离**:MySQL主从复制 +4. **添加监控告警**:Prometheus + Grafana +5. **完善API文档**:保持Swagger同步 + +### 11.4 长期演进建议(3-6月) + +1. **微服务拆分**:订单、支付、聊天独立服务 +2. **引入消息队列**:异步任务处理 +3. **数据库分库分表**:应对数据增长 +4. **容器编排**:Kubernetes部署 +5. **全链路监控**:分布式追踪(Jaeger) + +--- + +## 附录 + +### A. 技术债务清单 + +| 优先级 | 问题 | 影响 | 建议 | +|-------|------|------|------| +| P0 | 缺少单元测试 | 回归风险高 | 添加核心业务测试 | +| P1 | 缺少API限流 | 容易被刷 | 添加限流中间件 | +| P1 | 短信限流不足 | 验证码被刷 | 增加图形验证码 | +| P2 | 缺少缓存层 | 数据库压力大 | 添加Redis缓存 | +| P2 | 深度分页慢 | 用户体验差 | 游标分页 | +| P3 | 缺少监控 | 故障发现慢 | Prometheus | + +### B. 性能测试数据 + +**商品查询场景**(30并发60秒): +- QPS:4,106 +- P50延迟:5ms +- P95延迟:16ms +- 成功率:100% + +**管理后台场景**(待测试): +- 预期QPS:2,000+ +- 预期P95延迟:<50ms + +--- + +**报告编写人**:Claude Opus 4.8 +**审核状态**:已完成 +**下次更新**:根据性能测试结果动态调整 diff --git a/scripts/load_stress.go b/scripts/load_stress.go new file mode 100644 index 0000000..9d96a39 --- /dev/null +++ b/scripts/load_stress.go @@ -0,0 +1,832 @@ +package main + +import ( + "bytes" + "encoding/json" + "flag" + "fmt" + "io" + "math/rand" + "net/http" + "sort" + "sync" + "sync/atomic" + "time" +) + +// 改进的压力测试工具 - 支持真实认证、梯度压测、详细统计 + +type TestConfig struct { + BaseURL string + Scenario string + Duration time.Duration + Concurrency int + Gradual bool + WarmupUsers int +} + +type TestResult struct { + TotalRequests int64 + SuccessRequests int64 + FailedRequests int64 + Latencies []int64 // 存储所有延迟用于百分位计算 + Errors map[string]int64 + mu sync.Mutex +} + +type AuthPool struct { + tokens []string + mu sync.RWMutex +} + +type AdminToken struct { + token string + expiresAt time.Time + mu sync.RWMutex +} + +// 可用商品ID池 +type ListingIDPool struct { + ids []int + mu sync.RWMutex +} + +var globalListingIDs *ListingIDPool + +var ( + baseURL = flag.String("url", "http://localhost:8080", "API 基础地址") + concurrency = flag.Int("c", 50, "并发数") + duration = flag.Int("d", 60, "测试时长(秒)") + scenario = flag.String("s", "realistic", "测试场景: realistic(真实), admin(管理后台), listing_only(商品查询)") + gradual = flag.Bool("gradual", false, "启用梯度压测") + warmupUsers = flag.Int("warmup", 100, "预热用户数(生成token)") +) + +func main() { + flag.Parse() + + config := TestConfig{ + BaseURL: *baseURL, + Scenario: *scenario, + Duration: time.Duration(*duration) * time.Second, + Concurrency: *concurrency, + Gradual: *gradual, + WarmupUsers: *warmupUsers, + } + + fmt.Printf("=== 压力测试配置 ===\n") + fmt.Printf("目标地址: %s\n", config.BaseURL) + fmt.Printf("测试场景: %s\n", config.Scenario) + fmt.Printf("测试时长: %d 秒\n", *duration) + fmt.Printf("并发数: %d\n", config.Concurrency) + fmt.Printf("梯度压测: %v\n", config.Gradual) + fmt.Printf("==================\n\n") + + // 健康检查 + if !healthCheck(config.BaseURL) { + fmt.Println("❌ 后端服务未响应,退出测试") + return + } + + // 预加载可用商品ID + fmt.Println("⏳ 预加载可用商品ID...") + globalListingIDs = loadAvailableListings(config.BaseURL) + if globalListingIDs == nil || len(globalListingIDs.ids) == 0 { + fmt.Println("⚠️ 无法加载商品ID,将使用随机ID(可能导致404)") + } else { + fmt.Printf("✅ 成功加载 %d 个可用商品ID\n\n", len(globalListingIDs.ids)) + } + + // 根据场景初始化认证 + var authPool *AuthPool + var adminToken *AdminToken + + if config.Scenario == "realistic" || config.Scenario == "listing_only" { + fmt.Printf("⏳ 预热:生成 %d 个测试用户 token...\n", config.WarmupUsers) + authPool = initAuthPool(config.BaseURL, config.WarmupUsers) + if authPool == nil || len(authPool.tokens) == 0 { + fmt.Println("⚠️ 无法生成用户token,将使用匿名访问") + } else { + fmt.Printf("✅ 成功生成 %d 个用户 token\n\n", len(authPool.tokens)) + } + } + + if config.Scenario == "admin" { + fmt.Println("⏳ 获取管理员 token...") + adminToken = initAdminToken(config.BaseURL) + if adminToken == nil || adminToken.token == "" { + fmt.Println("❌ 无法获取管理员token,退出测试") + return + } + fmt.Println("✅ 成功获取管理员 token\n") + } + + // 执行压测 + var result *TestResult + if config.Gradual { + result = runGradualTest(config, authPool, adminToken) + } else { + result = runTest(config, authPool, adminToken) + } + + printResult(result, *duration) +} + +// ============================================ +// 健康检查和认证初始化 +// ============================================ + +func healthCheck(baseURL string) bool { + client := &http.Client{Timeout: 5 * time.Second} + resp, err := client.Get(baseURL + "/health") + if err != nil { + return false + } + defer resp.Body.Close() + return resp.StatusCode == 200 +} + +func initAuthPool(baseURL string, userCount int) *AuthPool { + pool := &AuthPool{tokens: make([]string, 0, userCount)} + client := &http.Client{Timeout: 10 * time.Second} + + successCount := 0 + for i := 1; i <= userCount; i++ { + phone := fmt.Sprintf("138%08d", i) + token := loginTestUser(client, baseURL, phone) + if token != "" { + pool.tokens = append(pool.tokens, token) + successCount++ + } + + // 每10个打印一次进度 + if i%10 == 0 || i == userCount { + fmt.Printf("\r 进度: %d/%d (%d 成功)", i, userCount, successCount) + } + } + fmt.Println() + + return pool +} + +func loginTestUser(client *http.Client, baseURL string, phone string) string { + // 直接登录,跳过发送验证码(避免限流) + // Mock provider的验证码固定为 123456,直接使用 + loginPayload := map[string]string{ + "phone": phone, + "code": "123456", + } + loginBody, _ := json.Marshal(loginPayload) + + resp, err := client.Post(baseURL+"/api/auth/sms/login", "application/json", bytes.NewBuffer(loginBody)) + if err != nil { + return "" + } + defer resp.Body.Close() + + // 如果验证码无效,发送一次验证码后重试 + if resp.StatusCode != 200 { + // 发送验证码 + sendPayload := map[string]string{"phone": phone} + sendBody, _ := json.Marshal(sendPayload) + sendResp, err := client.Post(baseURL+"/api/auth/sms/send", "application/json", bytes.NewBuffer(sendBody)) + if err != nil { + return "" + } + io.Copy(io.Discard, sendResp.Body) + sendResp.Body.Close() + + // 等待验证码写入Redis + time.Sleep(200 * time.Millisecond) + + // 重试登录 + resp2, err := client.Post(baseURL+"/api/auth/sms/login", "application/json", bytes.NewBuffer(loginBody)) + if err != nil { + return "" + } + defer resp2.Body.Close() + + if resp2.StatusCode != 200 { + return "" + } + + var result struct { + Code string `json:"code"` + Data struct { + AccessToken string `json:"access_token"` + } `json:"data"` + } + + if err := json.NewDecoder(resp2.Body).Decode(&result); err != nil { + return "" + } + + if result.Code == "ok" { + return result.Data.AccessToken + } + return "" + } + + var result struct { + Code string `json:"code"` // 修复:API返回字符串"ok" + Data struct { + AccessToken string `json:"access_token"` + } `json:"data"` + } + + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return "" + } + + if result.Code == "ok" { + return result.Data.AccessToken + } + + return "" +} + +func initAdminToken(baseURL string) *AdminToken { + client := &http.Client{Timeout: 10 * time.Second} + + // 获取验证码(为了获取captcha_id) + resp, err := client.Get(baseURL + "/api/admin/auth/captcha") + if err != nil { + fmt.Printf(" 错误: %v\n", err) + return nil + } + defer resp.Body.Close() + + var captchaResp struct { + Data struct { + CaptchaID string `json:"captcha_id"` + } `json:"data"` + } + + if err := json.NewDecoder(resp.Body).Decode(&captchaResp); err != nil { + return nil + } + + // 登录(使用默认管理员账号,验证码留空mock会通过) + loginPayload := map[string]string{ + "username": "admin", + "password": "admin123456", + "captcha_id": captchaResp.Data.CaptchaID, + "captcha": "1234", // mock模式会自动通过 + } + loginBody, _ := json.Marshal(loginPayload) + + resp, err = client.Post(baseURL+"/api/admin/auth/login", "application/json", bytes.NewBuffer(loginBody)) + if err != nil { + fmt.Printf(" 错误: %v\n", err) + return nil + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + fmt.Printf(" 登录失败: HTTP %d\n", resp.StatusCode) + return nil + } + + var loginResp struct { + Code string `json:"code"` + Data struct { + AccessToken string `json:"access_token"` + ExpiresIn int `json:"expires_in"` + } `json:"data"` + } + + if err := json.NewDecoder(resp.Body).Decode(&loginResp); err != nil { + return nil + } + + if loginResp.Code != "ok" { + fmt.Printf(" 登录失败: code=%s\n", loginResp.Code) + return nil + } + + return &AdminToken{ + token: loginResp.Data.AccessToken, + expiresAt: time.Now().Add(time.Duration(loginResp.Data.ExpiresIn) * time.Second), + } +} + +func (p *AuthPool) GetRandomToken() string { + if p == nil || len(p.tokens) == 0 { + return "" + } + p.mu.RLock() + defer p.mu.RUnlock() + return p.tokens[rand.Intn(len(p.tokens))] +} + +func (a *AdminToken) Get() string { + if a == nil { + return "" + } + a.mu.RLock() + defer a.mu.RUnlock() + return a.token +} + +func (p *ListingIDPool) GetRandomID() int { + if p == nil || len(p.ids) == 0 { + // 降级:返回随机ID + return rand.Intn(5000) + 1 + } + p.mu.RLock() + defer p.mu.RUnlock() + return p.ids[rand.Intn(len(p.ids))] +} + +func loadAvailableListings(baseURL string) *ListingIDPool { + client := &http.Client{Timeout: 10 * time.Second} + pool := &ListingIDPool{ids: make([]int, 0, 1000)} + + // 获取前1000个可用商品ID + for page := 1; page <= 10; page++ { + url := fmt.Sprintf("%s/api/listings?page=%d&page_size=100", baseURL, page) + resp, err := client.Get(url) + if err != nil { + break + } + + var result struct { + Code string `json:"code"` // 修复:API返回字符串"ok"而不是数字0 + Data struct { + Items []struct { + ID int `json:"id"` + } `json:"items"` + } `json:"data"` + } + + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + resp.Body.Close() + break + } + resp.Body.Close() + + if result.Code != "ok" || len(result.Data.Items) == 0 { + break + } + + for _, item := range result.Data.Items { + pool.ids = append(pool.ids, item.ID) + } + + // 如果不满100个,说明已经到最后一页 + if len(result.Data.Items) < 100 { + break + } + } + + return pool +} + +// ============================================ +// 测试执行 +// ============================================ + +func runTest(config TestConfig, authPool *AuthPool, adminToken *AdminToken) *TestResult { + result := &TestResult{ + Errors: make(map[string]int64), + Latencies: make([]int64, 0, 100000), + } + + var wg sync.WaitGroup + stopChan := make(chan struct{}) + + fmt.Printf("🚀 开始压测 (%d 并发, %v)...\n\n", config.Concurrency, config.Duration) + + // 启动统计goroutine + statsTicker := time.NewTicker(5 * time.Second) + go func() { + for { + select { + case <-statsTicker.C: + printProgress(result) + case <-stopChan: + statsTicker.Stop() + return + } + } + }() + + // 启动并发workers + for i := 0; i < config.Concurrency; i++ { + wg.Add(1) + go func(workerID int) { + defer wg.Done() + worker(workerID, config, result, authPool, adminToken, stopChan) + }(i) + } + + // 等待测试时长 + time.Sleep(config.Duration) + close(stopChan) + + wg.Wait() + + return result +} + +func runGradualTest(config TestConfig, authPool *AuthPool, adminToken *AdminToken) *TestResult { + profiles := []struct { + Duration time.Duration + Concurrency int + }{ + {30 * time.Second, 10}, // 预热 + {60 * time.Second, config.Concurrency / 4}, // 25%负载 + {60 * time.Second, config.Concurrency / 2}, // 50%负载 + {60 * time.Second, config.Concurrency}, // 100%负载 + {30 * time.Second, config.Concurrency * 2}, // 峰值负载 + } + + result := &TestResult{ + Errors: make(map[string]int64), + Latencies: make([]int64, 0, 200000), + } + + for i, profile := range profiles { + fmt.Printf("📊 阶段 %d/%d: %d 并发, 持续 %v\n", i+1, len(profiles), profile.Concurrency, profile.Duration) + + phaseConfig := config + phaseConfig.Duration = profile.Duration + phaseConfig.Concurrency = profile.Concurrency + + phaseResult := runTest(phaseConfig, authPool, adminToken) + + // 合并结果 + result.mu.Lock() + result.TotalRequests += phaseResult.TotalRequests + result.SuccessRequests += phaseResult.SuccessRequests + result.FailedRequests += phaseResult.FailedRequests + result.Latencies = append(result.Latencies, phaseResult.Latencies...) + for k, v := range phaseResult.Errors { + result.Errors[k] += v + } + result.mu.Unlock() + + if i < len(profiles)-1 { + fmt.Println("⏸ 冷却 10 秒...") + time.Sleep(10 * time.Second) + } + } + + return result +} + +func worker(id int, config TestConfig, result *TestResult, authPool *AuthPool, adminToken *AdminToken, stopChan chan struct{}) { + client := &http.Client{ + Timeout: 10 * time.Second, + } + + for { + select { + case <-stopChan: + return + default: + executeScenario(client, config, result, authPool, adminToken) + } + } +} + +func executeScenario(client *http.Client, config TestConfig, result *TestResult, authPool *AuthPool, adminToken *AdminToken) { + switch config.Scenario { + case "realistic": + executeRealisticScenario(client, config.BaseURL, result, authPool) + case "admin": + executeAdminScenario(client, config.BaseURL, result, adminToken) + case "listing_only": + executeListingOnlyScenario(client, config.BaseURL, result, authPool) + default: + testHealthCheck(client, config.BaseURL, result) + } +} + +// ============================================ +// 业务场景实现 +// ============================================ + +func executeRealisticScenario(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + r := rand.Intn(1000) + switch { + case r < 350: // 35% 查询商品列表 + testListListings(client, baseURL, result) + case r < 600: // 25% 查询商品详情 + testGetListingDetail(client, baseURL, result) + case r < 750: // 15% 查询我的订单 + testListMyOrders(client, baseURL, result, authPool) + case r < 850: // 10% 查询钱包余额 + testWalletBalance(client, baseURL, result, authPool) + case r < 920: // 7% 聊天列表 + testChatList(client, baseURL, result, authPool) + case r < 970: // 5% 查询订单详情 + testOrderDetail(client, baseURL, result, authPool) + case r < 990: // 2% 创建订单(需要认证) + testCreateOrder(client, baseURL, result, authPool) + default: // 1% 支付订单(需要认证) + testPayOrder(client, baseURL, result, authPool) + } +} + +func executeAdminScenario(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + r := rand.Intn(100) + switch { + case r < 25: // 25% 用户管理列表 + testAdminUsers(client, baseURL, result, adminToken) + case r < 45: // 20% 订单管理列表 + testAdminOrders(client, baseURL, result, adminToken) + case r < 60: // 15% 商品审核列表 + testAdminListings(client, baseURL, result, adminToken) + case r < 75: // 15% 钱包流水 + testAdminWalletLedger(client, baseURL, result, adminToken) + case r < 85: // 10% 申诉管理 + testAdminDisputes(client, baseURL, result, adminToken) + case r < 95: // 10% 审计日志 + testAdminAuditLogs(client, baseURL, result, adminToken) + default: // 5% 仪表盘 + testAdminDashboard(client, baseURL, result, adminToken) + } +} + +func executeListingOnlyScenario(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + r := rand.Intn(100) + if r < 70 { + testListListings(client, baseURL, result) + } else { + testGetListingDetail(client, baseURL, result) + } +} + +// ============================================ +// 具体测试函数 +// ============================================ + +func testHealthCheck(client *http.Client, baseURL string, result *TestResult) { + makeRequest(client, "GET", baseURL+"/health", "", nil, result, "health_check") +} + +func testListListings(client *http.Client, baseURL string, result *TestResult) { + page := rand.Intn(10) + 1 + pageSize := []int{10, 20, 50}[rand.Intn(3)] + url := fmt.Sprintf("%s/api/listings?page=%d&page_size=%d", baseURL, page, pageSize) + makeRequest(client, "GET", url, "", nil, result, "list_listings") +} + +func testGetListingDetail(client *http.Client, baseURL string, result *TestResult) { + listingID := globalListingIDs.GetRandomID() + url := fmt.Sprintf("%s/api/listings/%d", baseURL, listingID) + makeRequest(client, "GET", url, "", nil, result, "get_listing_detail") +} + +func testListMyOrders(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + token := authPool.GetRandomToken() + if token == "" { + return + } + page := rand.Intn(5) + 1 + url := fmt.Sprintf("%s/api/orders?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, token, nil, result, "list_my_orders") +} + +func testWalletBalance(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + token := authPool.GetRandomToken() + if token == "" { + return + } + makeAuthRequest(client, "GET", baseURL+"/api/wallet/balance", token, nil, result, "wallet_balance") +} + +func testChatList(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + token := authPool.GetRandomToken() + if token == "" { + return + } + makeAuthRequest(client, "GET", baseURL+"/api/chats?page=1&page_size=20", token, nil, result, "chat_list") +} + +func testOrderDetail(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + token := authPool.GetRandomToken() + if token == "" { + return + } + orderID := rand.Intn(3000) + 1 + url := fmt.Sprintf("%s/api/orders/%d", baseURL, orderID) + makeAuthRequest(client, "GET", url, token, nil, result, "order_detail") +} + +func testCreateOrder(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + token := authPool.GetRandomToken() + if token == "" { + return + } + listingID := globalListingIDs.GetRandomID() + payload := map[string]interface{}{ + "listing_id": listingID, + "estimated_duration_hours": 24, + } + body, _ := json.Marshal(payload) + makeAuthRequest(client, "POST", baseURL+"/api/orders", token, body, result, "create_order") +} + +func testPayOrder(client *http.Client, baseURL string, result *TestResult, authPool *AuthPool) { + token := authPool.GetRandomToken() + if token == "" { + return + } + orderID := rand.Intn(3000) + 1 + url := fmt.Sprintf("%s/api/orders/%d/start-payment", baseURL, orderID) + payload := map[string]interface{}{ + "provider": "mock", + } + body, _ := json.Marshal(payload) + makeAuthRequest(client, "POST", url, token, body, result, "pay_order") +} + +// 管理后台测试函数 +func testAdminUsers(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + page := rand.Intn(10) + 1 + url := fmt.Sprintf("%s/api/admin/users?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, adminToken.Get(), nil, result, "admin_users") +} + +func testAdminOrders(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + page := rand.Intn(10) + 1 + url := fmt.Sprintf("%s/api/admin/orders?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, adminToken.Get(), nil, result, "admin_orders") +} + +func testAdminListings(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + page := rand.Intn(10) + 1 + url := fmt.Sprintf("%s/api/admin/listings?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, adminToken.Get(), nil, result, "admin_listings") +} + +func testAdminWalletLedger(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + page := rand.Intn(20) + 1 + url := fmt.Sprintf("%s/api/admin/wallet/ledger?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, adminToken.Get(), nil, result, "admin_wallet_ledger") +} + +func testAdminDisputes(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + page := rand.Intn(5) + 1 + url := fmt.Sprintf("%s/api/admin/disputes?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, adminToken.Get(), nil, result, "admin_disputes") +} + +func testAdminAuditLogs(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + page := rand.Intn(20) + 1 + url := fmt.Sprintf("%s/api/admin/audit-logs?page=%d&page_size=20", baseURL, page) + makeAuthRequest(client, "GET", url, adminToken.Get(), nil, result, "admin_audit_logs") +} + +func testAdminDashboard(client *http.Client, baseURL string, result *TestResult, adminToken *AdminToken) { + makeAuthRequest(client, "GET", baseURL+"/api/admin/dashboard", adminToken.Get(), nil, result, "admin_dashboard") +} + +// ============================================ +// HTTP 请求辅助函数 +// ============================================ + +func makeRequest(client *http.Client, method, url, token string, body []byte, result *TestResult, apiName string) { + var req *http.Request + var err error + + if body != nil { + req, err = http.NewRequest(method, url, bytes.NewBuffer(body)) + } else { + req, err = http.NewRequest(method, url, nil) + } + + if err != nil { + recordError(result, apiName+"_req_error") + atomic.AddInt64(&result.TotalRequests, 1) + atomic.AddInt64(&result.FailedRequests, 1) + return + } + + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + if body != nil { + req.Header.Set("Content-Type", "application/json") + } + + start := time.Now() + resp, err := client.Do(req) + latency := time.Since(start).Milliseconds() + + atomic.AddInt64(&result.TotalRequests, 1) + recordLatency(result, latency) + + if err != nil { + atomic.AddInt64(&result.FailedRequests, 1) + recordError(result, apiName+"_error: "+err.Error()) + return + } + defer resp.Body.Close() + io.Copy(io.Discard, resp.Body) + + if resp.StatusCode >= 200 && resp.StatusCode < 300 { + atomic.AddInt64(&result.SuccessRequests, 1) + } else { + atomic.AddInt64(&result.FailedRequests, 1) + recordError(result, fmt.Sprintf("%s_status_%d", apiName, resp.StatusCode)) + } +} + +func makeAuthRequest(client *http.Client, method, url, token string, body []byte, result *TestResult, apiName string) { + makeRequest(client, method, url, token, body, result, apiName) +} + +func recordLatency(result *TestResult, latency int64) { + result.mu.Lock() + defer result.mu.Unlock() + result.Latencies = append(result.Latencies, latency) +} + +func recordError(result *TestResult, errMsg string) { + result.mu.Lock() + defer result.mu.Unlock() + result.Errors[errMsg]++ +} + +// ============================================ +// 结果统计和输出 +// ============================================ + +func printProgress(result *TestResult) { + total := atomic.LoadInt64(&result.TotalRequests) + success := atomic.LoadInt64(&result.SuccessRequests) + failed := atomic.LoadInt64(&result.FailedRequests) + + if total > 0 { + successRate := float64(success) / float64(total) * 100 + fmt.Printf(" 进行中: %d 请求 | 成功率: %.2f%% | 失败: %d\n", total, successRate, failed) + } +} + +func printResult(result *TestResult, durationSec int) { + fmt.Printf("\n\n=== 压力测试结果 ===\n") + fmt.Printf("总请求数: %d\n", result.TotalRequests) + fmt.Printf("成功请求: %d (%.2f%%)\n", + result.SuccessRequests, + float64(result.SuccessRequests)/float64(result.TotalRequests)*100) + fmt.Printf("失败请求: %d (%.2f%%)\n", + result.FailedRequests, + float64(result.FailedRequests)/float64(result.TotalRequests)*100) + + qps := float64(result.TotalRequests) / float64(durationSec) + fmt.Printf("\nQPS: %.2f\n", qps) + + if len(result.Latencies) > 0 { + result.mu.Lock() + latencies := make([]int64, len(result.Latencies)) + copy(latencies, result.Latencies) + result.mu.Unlock() + + sort.Slice(latencies, func(i, j int) bool { return latencies[i] < latencies[j] }) + + p50 := latencies[len(latencies)*50/100] + p95 := latencies[len(latencies)*95/100] + p99 := latencies[len(latencies)*99/100] + min := latencies[0] + max := latencies[len(latencies)-1] + + var sum int64 + for _, l := range latencies { + sum += l + } + avg := sum / int64(len(latencies)) + + fmt.Printf("\n延迟统计:\n") + fmt.Printf(" 最小: %d ms\n", min) + fmt.Printf(" P50: %d ms\n", p50) + fmt.Printf(" 平均: %d ms\n", avg) + fmt.Printf(" P95: %d ms\n", p95) + fmt.Printf(" P99: %d ms\n", p99) + fmt.Printf(" 最大: %d ms\n", max) + } + + if len(result.Errors) > 0 { + fmt.Printf("\n错误统计 (Top 10):\n") + type errorPair struct { + msg string + count int64 + } + var errors []errorPair + for msg, count := range result.Errors { + errors = append(errors, errorPair{msg, count}) + } + sort.Slice(errors, func(i, j int) bool { return errors[i].count > errors[j].count }) + + for i, e := range errors { + if i >= 10 { + break + } + fmt.Printf(" %s: %d 次\n", e.msg, e.count) + } + } + + fmt.Printf("==================\n") +} diff --git a/scripts/load_test_data.sql b/scripts/load_test_data.sql index 5c6348e..0c9a5ab 100644 --- a/scripts/load_test_data.sql +++ b/scripts/load_test_data.sql @@ -85,12 +85,32 @@ BEGIN DECLARE account_id_val BIGINT; DECLARE price_val DECIMAL(12,2); DECLARE deposit_val DECIMAL(12,2); + DECLARE verified_user_min BIGINT; + DECLARE verified_user_max BIGINT; + DECLARE user_pick BIGINT; + + SELECT MIN(id), MAX(id) + INTO verified_user_min, verified_user_max + FROM users + WHERE realname_status = 'verified'; + + IF verified_user_min IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的已实名用户'; + END IF; WHILE i <= batch_size DO - -- 随机选择一个用户作为号主(已实名用户) + -- 近似随机选择一个用户作为号主,避免 ORDER BY RAND() 全表排序。 + SET user_id_val = NULL; + SET user_pick = verified_user_min + FLOOR(RAND() * (verified_user_max - verified_user_min + 1)); SELECT id INTO user_id_val FROM users - WHERE realname_status = 'verified' - ORDER BY RAND() LIMIT 1; + WHERE id >= user_pick AND realname_status = 'verified' + ORDER BY id LIMIT 1; + + IF user_id_val IS NULL THEN + SELECT id INTO user_id_val FROM users + WHERE realname_status = 'verified' + ORDER BY id LIMIT 1; + END IF; -- 创建游戏账号 INSERT INTO game_accounts ( @@ -120,7 +140,7 @@ BEGIN ELSE '钻石' END, (i % 100) * 10000, - 'active' + 'published' ); SET account_id_val = LAST_INSERT_ID(); @@ -141,7 +161,7 @@ BEGIN CASE WHEN i % 20 = 0 THEN 'offline' WHEN i % 15 = 0 THEN 'draft' - ELSE 'active' + ELSE 'published' END, CASE WHEN i % 15 = 0 THEN 'pending' @@ -177,19 +197,67 @@ BEGIN DECLARE order_no_val VARCHAR(64); DECLARE rent_amount_val DECIMAL(12,2); DECLARE deposit_val DECIMAL(12,2); + DECLARE listing_min BIGINT; + DECLARE listing_max BIGINT; + DECLARE listing_pick BIGINT; + DECLARE verified_user_min BIGINT; + DECLARE verified_user_max BIGINT; + DECLARE user_pick BIGINT; + + SELECT MIN(id), MAX(id) + INTO listing_min, listing_max + FROM rental_listings + WHERE status IN ('published', 'active') AND review_status = 'approved'; + + SELECT MIN(id), MAX(id) + INTO verified_user_min, verified_user_max + FROM users + WHERE realname_status = 'verified'; + + IF listing_min IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的已上架商品'; + END IF; + IF verified_user_min IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的已实名用户'; + END IF; WHILE i <= batch_size DO - -- 随机选择一个上架的商品 + -- 近似随机选择一个上架商品,避免 ORDER BY RAND() 全表排序。 + SET listing_id_val = NULL; + SET listing_pick = listing_min + FLOOR(RAND() * (listing_max - listing_min + 1)); SELECT rl.id, rl.account_id, rl.owner_id, rl.price, rl.deposit_amount INTO listing_id_val, account_id_val, owner_id_val, rent_amount_val, deposit_val FROM rental_listings rl - WHERE rl.status = 'active' AND rl.review_status = 'approved' - ORDER BY RAND() LIMIT 1; + WHERE rl.id >= listing_pick + AND rl.status IN ('published', 'active') + AND rl.review_status = 'approved' + ORDER BY rl.id LIMIT 1; - -- 随机选择一个租客(不能是号主本人) + IF listing_id_val IS NULL THEN + SELECT rl.id, rl.account_id, rl.owner_id, rl.price, rl.deposit_amount + INTO listing_id_val, account_id_val, owner_id_val, rent_amount_val, deposit_val + FROM rental_listings rl + WHERE rl.status IN ('published', 'active') + AND rl.review_status = 'approved' + ORDER BY rl.id LIMIT 1; + END IF; + + -- 近似随机选择一个租客(不能是号主本人)。 + SET renter_id_val = NULL; + SET user_pick = verified_user_min + FLOOR(RAND() * (verified_user_max - verified_user_min + 1)); SELECT id INTO renter_id_val FROM users - WHERE id != owner_id_val AND realname_status = 'verified' - ORDER BY RAND() LIMIT 1; + WHERE id >= user_pick AND id != owner_id_val AND realname_status = 'verified' + ORDER BY id LIMIT 1; + + IF renter_id_val IS NULL THEN + SELECT id INTO renter_id_val FROM users + WHERE id != owner_id_val AND realname_status = 'verified' + ORDER BY id LIMIT 1; + END IF; + + IF renter_id_val IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的租客用户'; + END IF; SET order_no_val = CONCAT('ORD', DATE_FORMAT(NOW(), '%Y%m%d'), LPAD(i, 8, '0')); SET rent_amount_val = rent_amount_val * 24; -- 24小时租金 @@ -254,14 +322,43 @@ BEGIN DECLARE order_id_val BIGINT; DECLARE ledger_no_val VARCHAR(64); DECLARE amount_val DECIMAL(12,2); + DECLARE user_min BIGINT; + DECLARE user_max BIGINT; + DECLARE order_min BIGINT; + DECLARE order_max BIGINT; + DECLARE user_pick BIGINT; + DECLARE order_pick BIGINT; + + SELECT MIN(id), MAX(id) INTO user_min, user_max FROM users; + SELECT MIN(id), MAX(id) INTO order_min, order_max FROM rental_orders; + + IF user_min IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的用户'; + END IF; WHILE i <= batch_size DO - -- 随机选择用户 - SELECT id INTO user_id_val FROM users ORDER BY RAND() LIMIT 1; + -- 近似随机选择用户,避免 ORDER BY RAND() 全表排序。 + SET user_id_val = NULL; + SET user_pick = user_min + FLOOR(RAND() * (user_max - user_min + 1)); + SELECT id INTO user_id_val FROM users + WHERE id >= user_pick + ORDER BY id LIMIT 1; + + IF user_id_val IS NULL THEN + SELECT id INTO user_id_val FROM users ORDER BY id LIMIT 1; + END IF; -- 随机关联订单(50%概率) - IF RAND() > 0.5 THEN - SELECT id INTO order_id_val FROM rental_orders ORDER BY RAND() LIMIT 1; + IF RAND() > 0.5 AND order_min IS NOT NULL THEN + SET order_id_val = NULL; + SET order_pick = order_min + FLOOR(RAND() * (order_max - order_min + 1)); + SELECT id INTO order_id_val FROM rental_orders + WHERE id >= order_pick + ORDER BY id LIMIT 1; + + IF order_id_val IS NULL THEN + SELECT id INTO order_id_val FROM rental_orders ORDER BY id LIMIT 1; + END IF; ELSE SET order_id_val = NULL; END IF; @@ -314,14 +411,45 @@ BEGIN DECLARE i INT DEFAULT 1; DECLARE conv_id BIGINT; DECLARE user_id_val BIGINT; + DECLARE conv_min BIGINT; + DECLARE conv_max BIGINT; + DECLARE user_min BIGINT; + DECLARE user_max BIGINT; + DECLARE conv_pick BIGINT; + DECLARE user_pick BIGINT; + + SELECT MIN(id), MAX(id) INTO conv_min, conv_max FROM chat_conversations; + SELECT MIN(id), MAX(id) INTO user_min, user_max FROM users; + + IF conv_min IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的聊天会话'; + END IF; + IF user_min IS NULL THEN + SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '没有可用的用户'; + END IF; WHILE i <= batch_size DO - -- 随机选择一个会话 - SELECT id INTO conv_id FROM chat_conversations ORDER BY RAND() LIMIT 1; + -- 近似随机选择一个会话和发送者,避免 ORDER BY RAND() 全表排序。 + SET conv_id = NULL; + SET conv_pick = conv_min + FLOOR(RAND() * (conv_max - conv_min + 1)); + SELECT id INTO conv_id FROM chat_conversations + WHERE id >= conv_pick + ORDER BY id LIMIT 1; + + IF conv_id IS NULL THEN + SELECT id INTO conv_id FROM chat_conversations ORDER BY id LIMIT 1; + END IF; IF conv_id IS NOT NULL THEN - -- 随机选择发送者 - SELECT id INTO user_id_val FROM users ORDER BY RAND() LIMIT 1; + SET user_id_val = NULL; + SET user_pick = user_min + FLOOR(RAND() * (user_max - user_min + 1)); + SELECT id INTO user_id_val FROM users + WHERE id >= user_pick + ORDER BY id LIMIT 1; + + IF user_id_val IS NULL THEN + SELECT id INTO user_id_val FROM users ORDER BY id LIMIT 1; + END IF; INSERT INTO chat_messages ( conversation_id, sender_type, sender_id, sender_role, @@ -349,46 +477,8 @@ END$$ DELIMITER ; -- ============================================ --- 执行数据生成(根据需要调整数量) +-- 本文件只创建存储过程,不直接生成数据。 +-- 请使用 scripts/stress_test.sh data 按参数生成,避免意外重复造数。 -- ============================================ --- 生成10000个用户 -CALL generate_users(10000); - --- 生成50000个商品 -CALL generate_listings(50000); - --- 生成30000个订单 -CALL generate_orders(30000); - --- 生成100000条钱包流水 -CALL generate_wallet_ledger(100000); - --- 为前1000个订单创建会话 -INSERT INTO chat_conversations (order_id, type, title, status, last_message_at) -SELECT id, 'order_group', CONCAT('订单', order_no, '群聊'), 'active', created_at -FROM rental_orders -WHERE id <= 1000 -ON DUPLICATE KEY UPDATE order_id=order_id; - --- 生成50000条聊天消息 -CALL generate_chat_messages(50000); - --- ============================================ --- 查看数据统计 --- ============================================ -SELECT '用户数' as item, COUNT(*) as count FROM users -UNION ALL -SELECT '游戏账号数', COUNT(*) FROM game_accounts -UNION ALL -SELECT '商品数', COUNT(*) FROM rental_listings -UNION ALL -SELECT '订单数', COUNT(*) FROM rental_orders -UNION ALL -SELECT '钱包流水数', COUNT(*) FROM wallet_ledger -UNION ALL -SELECT '聊天会话数', COUNT(*) FROM chat_conversations -UNION ALL -SELECT '聊天消息数', COUNT(*) FROM chat_messages; - SET FOREIGN_KEY_CHECKS = 1; diff --git a/scripts/load_test_data_optimized.sql b/scripts/load_test_data_optimized.sql new file mode 100644 index 0000000..d0f4270 --- /dev/null +++ b/scripts/load_test_data_optimized.sql @@ -0,0 +1,475 @@ +-- ============================================ +-- 优化的压力测试数据生成脚本 +-- 使用批量生成 + 临时表,避免循环中的随机查询 +-- ============================================ + +SET NAMES utf8mb4; +SET FOREIGN_KEY_CHECKS = 0; + +-- ============================================ +-- 1. 批量生成用户数据 +-- ============================================ +DROP PROCEDURE IF EXISTS generate_users_batch; +DELIMITER $$ +CREATE PROCEDURE generate_users_batch(IN batch_size INT) +BEGIN + DECLARE batch_limit INT DEFAULT 1000; + DECLARE batches INT; + DECLARE current_batch INT DEFAULT 0; + DECLARE batch_start INT; + DECLARE batch_end INT; + + SET batches = CEIL(batch_size / batch_limit); + + WHILE current_batch < batches DO + SET batch_start = current_batch * batch_limit + 1; + SET batch_end = LEAST((current_batch + 1) * batch_limit, batch_size); + + -- 使用 INSERT ... SELECT 批量生成 + INSERT INTO users (phone, nickname, realname_status, risk_status, credit_score, status, created_at) + SELECT + CONCAT('138', LPAD(seq, 8, '0')) as phone, + CONCAT('测试用户', seq) as nickname, + CASE WHEN seq % 10 = 0 THEN 'unverified' ELSE 'verified' END as realname_status, + CASE WHEN seq % 100 = 0 THEN 'frozen' ELSE 'normal' END as risk_status, + 80 + (seq % 20) as credit_score, + 'active' as status, + DATE_SUB(NOW(), INTERVAL (seq % 365) DAY) as created_at + FROM ( + SELECT @row := @row + 1 AS seq + FROM + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t1, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t2, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t3, + (SELECT @row := batch_start - 1) r + LIMIT batch_end - batch_start + 1 + ) seqs + ON DUPLICATE KEY UPDATE id=id; + + SET current_batch = current_batch + 1; + COMMIT; + END WHILE; + + -- 批量生成实名记录(为已实名用户) + INSERT INTO user_realname (user_id, provider, status, masked_name, masked_id_no, verified_at) + SELECT + u.id, + 'mock', + 'success', + CONCAT('张*', CHAR(65 + (u.id % 26))), + CONCAT('3301**********', LPAD(u.id % 10000, 4, '0')), + DATE_SUB(NOW(), INTERVAL (u.id % 300) DAY) + FROM users u + WHERE u.realname_status = 'verified' + AND NOT EXISTS (SELECT 1 FROM user_realname WHERE user_id = u.id); + + -- 批量生成钱包账户 + INSERT INTO wallet_accounts (user_id, available_balance, frozen_balance, status) + SELECT + u.id, + (u.id % 10) * 100.00, + (u.id % 5) * 50.00, + 'active' + FROM users u + WHERE NOT EXISTS (SELECT 1 FROM wallet_accounts WHERE user_id = u.id); + + COMMIT; +END$$ +DELIMITER ; + +-- ============================================ +-- 2. 批量生成游戏账号和商品 +-- ============================================ +DROP PROCEDURE IF EXISTS generate_listings_batch; +DELIMITER $$ +CREATE PROCEDURE generate_listings_batch(IN batch_size INT) +BEGIN + DECLARE batch_limit INT DEFAULT 1000; + DECLARE batches INT; + DECLARE current_batch INT DEFAULT 0; + + -- 创建临时表存储已实名用户ID + DROP TEMPORARY TABLE IF EXISTS tmp_verified_users; + CREATE TEMPORARY TABLE tmp_verified_users ( + id BIGINT PRIMARY KEY, + row_num INT + ); + + INSERT INTO tmp_verified_users (id, row_num) + SELECT id, (@rn := @rn + 1) as row_num + FROM users, (SELECT @rn := 0) init + WHERE realname_status = 'verified' + ORDER BY id; + + SET batches = CEIL(batch_size / batch_limit); + + WHILE current_batch < batches DO + -- 批量生成游戏账号 + INSERT INTO game_accounts ( + owner_id, game_name, server_region, login_platform, + title, description, rank_level, haf_coin_amount, status + ) + SELECT + u.id as owner_id, + 'delta_force', + CASE (seq % 4) + WHEN 0 THEN '亚服' + WHEN 1 THEN '美服' + WHEN 2 THEN '欧服' + ELSE '国服' + END, + CASE (seq % 3) + WHEN 0 THEN 'Steam' + WHEN 1 THEN 'Epic' + ELSE 'WeGame' + END, + CONCAT('账号', seq, ' 高分段'), + CONCAT('这是一个测试账号,编号', seq), + CASE (seq % 5) + WHEN 0 THEN '青铜' + WHEN 1 THEN '白银' + WHEN 2 THEN '黄金' + WHEN 3 THEN '铂金' + ELSE '钻石' + END, + (seq % 100) * 10000, + 'published' + FROM ( + SELECT @row2 := @row2 + 1 AS seq + FROM + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t1, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t2, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t3, + (SELECT @row2 := current_batch * batch_limit) r + LIMIT batch_limit + ) seqs + INNER JOIN tmp_verified_users u ON u.row_num = (seq % (SELECT COUNT(*) FROM tmp_verified_users)) + 1 + WHERE seq <= batch_size; + + -- 批量生成租号商品 + INSERT INTO rental_listings ( + account_id, owner_id, price, deposit_amount, + in_transaction, status, review_status, published_at + ) + SELECT + ga.id as account_id, + ga.owner_id, + 5.00 + ((ga.id % 20) * 0.5) as price, + 100.00 + ((ga.id % 10) * 50.00) as deposit_amount, + CASE WHEN ga.id % 10 = 0 THEN 1 ELSE 0 END as in_transaction, + CASE + WHEN ga.id % 20 = 0 THEN 'offline' + WHEN ga.id % 15 = 0 THEN 'draft' + ELSE 'published' + END as status, + CASE + WHEN ga.id % 15 = 0 THEN 'pending' + WHEN ga.id % 30 = 0 THEN 'rejected' + ELSE 'approved' + END as review_status, + DATE_SUB(NOW(), INTERVAL (ga.id % 90) DAY) as published_at + FROM game_accounts ga + WHERE ga.id > (SELECT COALESCE(MAX(account_id), 0) FROM rental_listings) + LIMIT batch_limit; + + SET current_batch = current_batch + 1; + COMMIT; + END WHILE; + + DROP TEMPORARY TABLE IF EXISTS tmp_verified_users; +END$$ +DELIMITER ; + +-- ============================================ +-- 3. 批量生成订单数据(优化版) +-- ============================================ +DROP PROCEDURE IF EXISTS generate_orders_batch; +DELIMITER $$ +CREATE PROCEDURE generate_orders_batch(IN batch_size INT) +BEGIN + DECLARE batch_limit INT DEFAULT 1000; + DECLARE batches INT; + DECLARE current_batch INT DEFAULT 0; + + -- 创建临时表:可用商品 + DROP TEMPORARY TABLE IF EXISTS tmp_available_listings; + CREATE TEMPORARY TABLE tmp_available_listings ( + id BIGINT PRIMARY KEY, + account_id BIGINT, + owner_id BIGINT, + price DECIMAL(12,2), + deposit_amount DECIMAL(12,2), + row_num INT + ); + + INSERT INTO tmp_available_listings (id, account_id, owner_id, price, deposit_amount, row_num) + SELECT id, account_id, owner_id, price, deposit_amount, (@rn := @rn + 1) + FROM rental_listings, (SELECT @rn := 0) init + WHERE status IN ('published', 'active') AND review_status = 'approved' + ORDER BY id; + + -- 创建临时表:已实名用户 + DROP TEMPORARY TABLE IF EXISTS tmp_verified_users; + CREATE TEMPORARY TABLE tmp_verified_users ( + id BIGINT PRIMARY KEY, + row_num INT + ); + + INSERT INTO tmp_verified_users (id, row_num) + SELECT id, (@rn2 := @rn2 + 1) + FROM users, (SELECT @rn2 := 0) init + WHERE realname_status = 'verified' + ORDER BY id; + + SET batches = CEIL(batch_size / batch_limit); + + WHILE current_batch < batches DO + INSERT INTO rental_orders ( + order_no, listing_id, account_id, owner_id, renter_id, + estimated_duration_hours, rent_amount, owner_rent_amount, + deposit_amount, platform_fee, status, handoff_status, + settlement_status, rented_at, created_at + ) + SELECT + CONCAT('ORD', DATE_FORMAT(NOW(), '%Y%m%d'), LPAD(seq, 8, '0')) as order_no, + l.id as listing_id, + l.account_id, + l.owner_id, + r.id as renter_id, + 24 as estimated_duration_hours, + l.price * 24 as rent_amount, + l.price * 24 * 0.95 as owner_rent_amount, + l.deposit_amount, + l.price * 24 * 0.05 as platform_fee, + CASE (seq % 10) + WHEN 0 THEN 'pending_payment' + WHEN 1 THEN 'cancelled' + WHEN 2 THEN 'closed' + ELSE 'completed' + END as status, + CASE (seq % 10) + WHEN 0 THEN 'none' + WHEN 1 THEN 'none' + WHEN 2 THEN 'owner_delivered' + ELSE 'owner_received' + END as handoff_status, + CASE (seq % 10) + WHEN 0 THEN 'unsettled' + WHEN 1 THEN 'unsettled' + ELSE 'settled' + END as settlement_status, + DATE_SUB(NOW(), INTERVAL (seq % 60) DAY) as rented_at, + DATE_SUB(NOW(), INTERVAL (seq % 60) DAY) as created_at + FROM ( + SELECT @row3 := @row3 + 1 AS seq + FROM + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t1, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t2, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t3, + (SELECT @row3 := current_batch * batch_limit) r + LIMIT batch_limit + ) seqs + INNER JOIN tmp_available_listings l ON l.row_num = (seq % (SELECT COUNT(*) FROM tmp_available_listings)) + 1 + INNER JOIN tmp_verified_users r ON r.row_num = (seq % (SELECT COUNT(*) FROM tmp_verified_users)) + 1 + WHERE seq <= batch_size AND r.id != l.owner_id + LIMIT batch_limit; + + SET current_batch = current_batch + 1; + COMMIT; + END WHILE; + + DROP TEMPORARY TABLE IF EXISTS tmp_available_listings; + DROP TEMPORARY TABLE IF EXISTS tmp_verified_users; +END$$ +DELIMITER ; + +-- ============================================ +-- 4. 批量生成钱包流水 +-- ============================================ +DROP PROCEDURE IF EXISTS generate_wallet_ledger_batch; +DELIMITER $$ +CREATE PROCEDURE generate_wallet_ledger_batch(IN batch_size INT) +BEGIN + DECLARE batch_limit INT DEFAULT 1000; + DECLARE batches INT; + DECLARE current_batch INT DEFAULT 0; + + -- 创建临时表:用户列表 + DROP TEMPORARY TABLE IF EXISTS tmp_users; + CREATE TEMPORARY TABLE tmp_users ( + id BIGINT PRIMARY KEY, + row_num INT + ); + + INSERT INTO tmp_users (id, row_num) + SELECT id, (@rn := @rn + 1) + FROM users, (SELECT @rn := 0) init + ORDER BY id; + + -- 创建临时表:订单列表 + DROP TEMPORARY TABLE IF EXISTS tmp_orders; + CREATE TEMPORARY TABLE tmp_orders ( + id BIGINT PRIMARY KEY, + row_num INT + ); + + INSERT INTO tmp_orders (id, row_num) + SELECT id, (@rn2 := @rn2 + 1) + FROM rental_orders, (SELECT @rn2 := 0) init + ORDER BY id; + + SET batches = CEIL(batch_size / batch_limit); + + WHILE current_batch < batches DO + INSERT INTO wallet_ledger ( + ledger_no, user_id, order_id, direction, amount, + balance_after, balance_type, biz_type, biz_no, remark, created_at + ) + SELECT + CONCAT('LDG', DATE_FORMAT(NOW(), '%Y%m%d%H%i%s'), LPAD(seq, 6, '0')) as ledger_no, + u.id as user_id, + IF(seq % 2 = 0, o.id, NULL) as order_id, + CASE WHEN seq % 2 = 0 THEN 'in' ELSE 'out' END as direction, + (seq % 500) + (seq * 0.01) as amount, + 1000.00 + (seq % 1000) as balance_after, + CASE WHEN seq % 5 = 0 THEN 'frozen' ELSE 'available' END as balance_type, + CASE (seq % 6) + WHEN 0 THEN 'rent_payment' + WHEN 1 THEN 'deposit_freeze' + WHEN 2 THEN 'settlement' + WHEN 3 THEN 'refund' + WHEN 4 THEN 'recharge' + ELSE 'withdraw' + END as biz_type, + CONCAT('BIZ', LPAD(seq, 10, '0')) as biz_no, + CONCAT('测试流水', seq) as remark, + DATE_SUB(NOW(), INTERVAL (seq % 180) DAY) as created_at + FROM ( + SELECT @row4 := @row4 + 1 AS seq + FROM + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t1, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t2, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t3, + (SELECT @row4 := current_batch * batch_limit) r + LIMIT batch_limit + ) seqs + INNER JOIN tmp_users u ON u.row_num = (seq % (SELECT COUNT(*) FROM tmp_users)) + 1 + LEFT JOIN tmp_orders o ON o.row_num = (seq % (SELECT COUNT(*) FROM tmp_orders)) + 1 + WHERE seq <= batch_size + LIMIT batch_limit; + + SET current_batch = current_batch + 1; + COMMIT; + END WHILE; + + DROP TEMPORARY TABLE IF EXISTS tmp_users; + DROP TEMPORARY TABLE IF EXISTS tmp_orders; +END$$ +DELIMITER ; + +-- ============================================ +-- 5. 批量生成聊天会话和消息 +-- ============================================ +DROP PROCEDURE IF EXISTS generate_chat_data_batch; +DELIMITER $$ +CREATE PROCEDURE generate_chat_data_batch(IN message_count INT) +BEGIN + -- 先生成聊天会话(基于订单) + INSERT INTO chat_conversations (order_id, type, title, status, last_message_at) + SELECT + ro.id, + 'order_group', + CONCAT('订单', ro.order_no, '群聊'), + 'active', + ro.created_at + FROM rental_orders ro + WHERE NOT EXISTS (SELECT 1 FROM chat_conversations WHERE order_id = ro.id) + LIMIT 1000 + ON DUPLICATE KEY UPDATE order_id=order_id; + + -- 批量生成聊天消息 + DECLARE batch_limit INT DEFAULT 1000; + DECLARE batches INT; + DECLARE current_batch INT DEFAULT 0; + + DROP TEMPORARY TABLE IF EXISTS tmp_conversations; + CREATE TEMPORARY TABLE tmp_conversations ( + id BIGINT PRIMARY KEY, + row_num INT + ); + + INSERT INTO tmp_conversations (id, row_num) + SELECT id, (@rn := @rn + 1) + FROM chat_conversations, (SELECT @rn := 0) init + ORDER BY id; + + DROP TEMPORARY TABLE IF EXISTS tmp_users; + CREATE TEMPORARY TABLE tmp_users ( + id BIGINT PRIMARY KEY, + row_num INT + ); + + INSERT INTO tmp_users (id, row_num) + SELECT id, (@rn2 := @rn2 + 1) + FROM users, (SELECT @rn2 := 0) init + ORDER BY id; + + SET batches = CEIL(message_count / batch_limit); + + WHILE current_batch < batches DO + INSERT INTO chat_messages ( + conversation_id, sender_type, sender_id, sender_role, + content_type, content, created_at + ) + SELECT + c.id as conversation_id, + 'user' as sender_type, + u.id as sender_id, + CASE WHEN seq % 2 = 0 THEN 'owner' ELSE 'renter' END as sender_role, + 'text' as content_type, + CONCAT('这是测试消息', seq, ',内容随机生成用于压力测试') as content, + DATE_SUB(NOW(), INTERVAL (seq % 30) DAY) as created_at + FROM ( + SELECT @row5 := @row5 + 1 AS seq + FROM + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t1, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t2, + (SELECT 0 UNION SELECT 1 UNION SELECT 2 UNION SELECT 3 UNION SELECT 4 + UNION SELECT 5 UNION SELECT 6 UNION SELECT 7 UNION SELECT 8 UNION SELECT 9) t3, + (SELECT @row5 := current_batch * batch_limit) r + LIMIT batch_limit + ) seqs + INNER JOIN tmp_conversations c ON c.row_num = (seq % (SELECT COUNT(*) FROM tmp_conversations)) + 1 + INNER JOIN tmp_users u ON u.row_num = (seq % (SELECT COUNT(*) FROM tmp_users)) + 1 + WHERE seq <= message_count + LIMIT batch_limit; + + SET current_batch = current_batch + 1; + COMMIT; + END WHILE; + + DROP TEMPORARY TABLE IF EXISTS tmp_conversations; + DROP TEMPORARY TABLE IF EXISTS tmp_users; +END$$ +DELIMITER ; + +-- ============================================ +-- 本文件只创建存储过程,不直接生成数据。 +-- 请使用 scripts/stress_test.sh data 按参数生成 +-- ============================================ + +SET FOREIGN_KEY_CHECKS = 1; diff --git a/scripts/stress_test.go b/scripts/stress_test.go deleted file mode 100644 index 9c19063..0000000 --- a/scripts/stress_test.go +++ /dev/null @@ -1,407 +0,0 @@ -package main - -import ( - "bytes" - "encoding/json" - "flag" - "fmt" - "math/rand" - "net/http" - "sync" - "sync/atomic" - "time" -) - -// 压力测试工具 - 模拟实际业务场景 - -type TestConfig struct { - BaseURL string - Concurrency int - Duration time.Duration - Scenario string -} - -type TestResult struct { - TotalRequests int64 - SuccessRequests int64 - FailedRequests int64 - TotalLatency int64 // 毫秒 - MinLatency int64 - MaxLatency int64 - Errors map[string]int64 -} - -var ( - baseURL = flag.String("url", "http://localhost:8080", "API 基础地址") - concurrency = flag.Int("c", 10, "并发数") - duration = flag.Int("d", 60, "测试时长(秒)") - scenario = flag.String("s", "mixed", "测试场景: list_listings, create_order, chat, wallet, mixed") -) - -func main() { - flag.Parse() - - config := TestConfig{ - BaseURL: *baseURL, - Concurrency: *concurrency, - Duration: time.Duration(*duration) * time.Second, - Scenario: *scenario, - } - - fmt.Printf("=== 压力测试配置 ===\n") - fmt.Printf("目标地址: %s\n", config.BaseURL) - fmt.Printf("并发数: %d\n", config.Concurrency) - fmt.Printf("测试时长: %d 秒\n", *duration) - fmt.Printf("测试场景: %s\n", config.Scenario) - fmt.Printf("==================\n\n") - - result := runTest(config) - printResult(result) -} - -func runTest(config TestConfig) *TestResult { - result := &TestResult{ - Errors: make(map[string]int64), - MinLatency: int64(^uint64(0) >> 1), // Max int64 - } - - var wg sync.WaitGroup - stopChan := make(chan struct{}) - - // 启动并发workers - for i := 0; i < config.Concurrency; i++ { - wg.Add(1) - go func(workerID int) { - defer wg.Done() - worker(workerID, config, result, stopChan) - }(i) - } - - // 等待测试时长 - time.Sleep(config.Duration) - close(stopChan) - - wg.Wait() - - return result -} - -func worker(id int, config TestConfig, result *TestResult, stopChan chan struct{}) { - client := &http.Client{ - Timeout: 10 * time.Second, - } - - for { - select { - case <-stopChan: - return - default: - executeScenario(client, config, result) - } - } -} - -func executeScenario(client *http.Client, config TestConfig, result *TestResult) { - switch config.Scenario { - case "list_listings": - testListListings(client, config.BaseURL, result) - case "create_order": - testCreateOrder(client, config.BaseURL, result) - case "chat": - testChatMessages(client, config.BaseURL, result) - case "wallet": - testWalletLedger(client, config.BaseURL, result) - case "mixed": - // 混合场景:按实际业务比例分配 - r := rand.Intn(100) - switch { - case r < 40: // 40% 查询商品列表 - testListListings(client, config.BaseURL, result) - case r < 60: // 20% 查询订单 - testListOrders(client, config.BaseURL, result) - case r < 75: // 15% 查询钱包流水 - testWalletLedger(client, config.BaseURL, result) - case r < 85: // 10% 聊天消息 - testChatMessages(client, config.BaseURL, result) - case r < 95: // 10% 创建订单 - testCreateOrder(client, config.BaseURL, result) - default: // 5% 支付 - testPayOrder(client, config.BaseURL, result) - } - default: - testHealthCheck(client, config.BaseURL, result) - } -} - -// ============================================ -// 测试场景实现 -// ============================================ - -func testHealthCheck(client *http.Client, baseURL string, result *TestResult) { - start := time.Now() - resp, err := client.Get(baseURL + "/health") - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "health_check_error: "+err.Error()) - return - } - defer resp.Body.Close() - - if resp.StatusCode == 200 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("health_check_status_%d", resp.StatusCode)) - } -} - -func testListListings(client *http.Client, baseURL string, result *TestResult) { - // 模拟不同的查询条件 - page := rand.Intn(10) + 1 - pageSize := []int{10, 20, 50}[rand.Intn(3)] - url := fmt.Sprintf("%s/api/listings?page=%d&page_size=%d", baseURL, page, pageSize) - - start := time.Now() - resp, err := client.Get(url) - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "list_listings_error: "+err.Error()) - return - } - defer resp.Body.Close() - - if resp.StatusCode == 200 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("list_listings_status_%d", resp.StatusCode)) - } -} - -func testListOrders(client *http.Client, baseURL string, result *TestResult) { - // 需要登录token,这里模拟匿名访问(会返回401) - page := rand.Intn(5) + 1 - url := fmt.Sprintf("%s/api/orders?page=%d&page_size=20", baseURL, page) - - start := time.Now() - resp, err := client.Get(url) - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "list_orders_error: "+err.Error()) - return - } - defer resp.Body.Close() - - // 401是预期的(未登录) - if resp.StatusCode == 200 || resp.StatusCode == 401 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("list_orders_status_%d", resp.StatusCode)) - } -} - -func testCreateOrder(client *http.Client, baseURL string, result *TestResult) { - // 模拟创建订单(需要登录,会返回401) - listingID := rand.Intn(1000) + 1 - payload := map[string]interface{}{ - "listing_id": listingID, - "estimated_duration_hours": 24, - } - - body, _ := json.Marshal(payload) - start := time.Now() - resp, err := client.Post( - baseURL+"/api/orders", - "application/json", - bytes.NewBuffer(body), - ) - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "create_order_error: "+err.Error()) - return - } - defer resp.Body.Close() - - // 401是预期的(未登录) - if resp.StatusCode == 200 || resp.StatusCode == 401 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("create_order_status_%d", resp.StatusCode)) - } -} - -func testPayOrder(client *http.Client, baseURL string, result *TestResult) { - orderID := rand.Intn(1000) + 1 - url := fmt.Sprintf("%s/api/orders/%d/pay", baseURL, orderID) - - payload := map[string]interface{}{ - "provider": "mock", - } - - body, _ := json.Marshal(payload) - start := time.Now() - resp, err := client.Post(url, "application/json", bytes.NewBuffer(body)) - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "pay_order_error: "+err.Error()) - return - } - defer resp.Body.Close() - - // 401是预期的(未登录) - if resp.StatusCode == 200 || resp.StatusCode == 401 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("pay_order_status_%d", resp.StatusCode)) - } -} - -func testWalletLedger(client *http.Client, baseURL string, result *TestResult) { - page := rand.Intn(10) + 1 - url := fmt.Sprintf("%s/api/wallet/ledger?page=%d&page_size=20", baseURL, page) - - start := time.Now() - resp, err := client.Get(url) - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "wallet_ledger_error: "+err.Error()) - return - } - defer resp.Body.Close() - - // 401是预期的(未登录) - if resp.StatusCode == 200 || resp.StatusCode == 401 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("wallet_ledger_status_%d", resp.StatusCode)) - } -} - -func testChatMessages(client *http.Client, baseURL string, result *TestResult) { - conversationID := rand.Intn(100) + 1 - url := fmt.Sprintf("%s/api/chats/%d/messages?page=1&page_size=50", baseURL, conversationID) - - start := time.Now() - resp, err := client.Get(url) - latency := time.Since(start).Milliseconds() - - atomic.AddInt64(&result.TotalRequests, 1) - updateLatency(result, latency) - - if err != nil { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, "chat_messages_error: "+err.Error()) - return - } - defer resp.Body.Close() - - // 401是预期的(未登录) - if resp.StatusCode == 200 || resp.StatusCode == 401 { - atomic.AddInt64(&result.SuccessRequests, 1) - } else { - atomic.AddInt64(&result.FailedRequests, 1) - recordError(result, fmt.Sprintf("chat_messages_status_%d", resp.StatusCode)) - } -} - -// ============================================ -// 辅助函数 -// ============================================ - -func updateLatency(result *TestResult, latency int64) { - atomic.AddInt64(&result.TotalLatency, latency) - - // 更新最小延迟 - for { - current := atomic.LoadInt64(&result.MinLatency) - if latency >= current { - break - } - if atomic.CompareAndSwapInt64(&result.MinLatency, current, latency) { - break - } - } - - // 更新最大延迟 - for { - current := atomic.LoadInt64(&result.MaxLatency) - if latency <= current { - break - } - if atomic.CompareAndSwapInt64(&result.MaxLatency, current, latency) { - break - } - } -} - -var errorMutex sync.Mutex - -func recordError(result *TestResult, errMsg string) { - errorMutex.Lock() - defer errorMutex.Unlock() - result.Errors[errMsg]++ -} - -func printResult(result *TestResult) { - fmt.Printf("\n=== 压力测试结果 ===\n") - fmt.Printf("总请求数: %d\n", result.TotalRequests) - fmt.Printf("成功请求: %d (%.2f%%)\n", - result.SuccessRequests, - float64(result.SuccessRequests)/float64(result.TotalRequests)*100) - fmt.Printf("失败请求: %d (%.2f%%)\n", - result.FailedRequests, - float64(result.FailedRequests)/float64(result.TotalRequests)*100) - - if result.TotalRequests > 0 { - avgLatency := result.TotalLatency / result.TotalRequests - fmt.Printf("\n延迟统计:\n") - fmt.Printf(" 最小延迟: %d ms\n", result.MinLatency) - fmt.Printf(" 平均延迟: %d ms\n", avgLatency) - fmt.Printf(" 最大延迟: %d ms\n", result.MaxLatency) - } - - if len(result.Errors) > 0 { - fmt.Printf("\n错误统计:\n") - for err, count := range result.Errors { - fmt.Printf(" %s: %d 次\n", err, count) - } - } - - qps := float64(result.TotalRequests) / float64(*duration) - fmt.Printf("\nQPS: %.2f\n", qps) - fmt.Printf("==================\n") -} diff --git a/scripts/stress_test.sh b/scripts/stress_test.sh index 84af2f6..ab064b3 100755 --- a/scripts/stress_test.sh +++ b/scripts/stress_test.sh @@ -27,41 +27,58 @@ function log_error() { function show_usage() { cat << EOF -压力测试脚本 +压力测试脚本(优化版) 用法: $0 [command] [options] 命令: - data 生成测试数据 - test 执行压力测试 + data 生成测试数据(使用优化的批量生成) + test 执行压力测试(支持真实认证) monitor 监控系统性能 clean 清理测试数据 report 生成压测报告 all 执行完整流程(生成数据 + 压测 + 报告) 选项: - -u, --users NUM 生成用户数量(默认: 10000) - -l, --listings NUM 生成商品数量(默认: 50000) - -o, --orders NUM 生成订单数量(默认: 30000) - -c, --concurrency NUM 并发数(默认: 100) - -d, --duration SEC 测试时长/秒(默认: 60) - -s, --scenario NAME 测试场景: list_listings, create_order, wallet, chat, mixed(默认: mixed) - --url URL 后端地址(默认: http://localhost:8080) - -h, --help 显示帮助信息 + 数据生成: + -u, --users NUM 生成用户数量(默认: 1000) + -l, --listings NUM 生成商品数量(默认: 5000) + -o, --orders NUM 生成订单数量(默认: 3000) + --ledger NUM 生成钱包流水数量(默认: 10000) + --chat-messages NUM 生成聊天消息数量(默认: 5000) + --use-optimized 使用优化的数据生成脚本(推荐) + + 压力测试: + -c, --concurrency NUM 并发数(默认: 50) + -d, --duration SEC 测试时长/秒(默认: 60) + -s, --scenario NAME 测试场景: + - realistic: 真实业务场景(默认) + - admin: 管理后台场景 + - listing_only: 仅商品查询 + --gradual 启用梯度压测 + --warmup NUM 预热用户数(默认: 100) + --use-improved 使用改进的压测工具(支持真实认证) + --url URL 后端地址(默认: http://localhost:8080) + + 通用: + -h, --help 显示帮助信息 示例: - # 生成测试数据 - $0 data + # 小规模数据生成(使用优化脚本) + $0 data --use-optimized - # 执行混合场景压测(100并发,持续60秒) - $0 test -c 100 -d 60 -s mixed + # 真实场景压测(100个用户token,50并发,持续60秒) + $0 test --use-improved -c 50 -d 60 -s realistic --warmup 100 - # 执行商品列表查询压测 - $0 test -s list_listings -c 200 -d 120 + # 管理后台压测 + $0 test --use-improved -c 30 -d 120 -s admin - # 执行完整流程 - $0 all + # 梯度压测(逐步增加负载) + $0 test --use-improved --gradual -d 60 + + # 执行完整流程(优化版) + $0 all --use-optimized --use-improved # 清理测试数据 $0 clean @@ -70,14 +87,19 @@ EOF } # 默认参数 -USERS=10000 -LISTINGS=50000 -ORDERS=30000 -LEDGER=100000 -CONCURRENCY=100 +USERS=1000 +LISTINGS=5000 +ORDERS=3000 +LEDGER=10000 +CHAT_MESSAGES=5000 +CONCURRENCY=50 DURATION=60 -SCENARIO="mixed" +SCENARIO="realistic" BASE_URL="http://localhost:8080" +USE_OPTIMIZED=false +USE_IMPROVED=true +GRADUAL=false +WARMUP=100 # 解析命令行参数 COMMAND="" @@ -99,6 +121,14 @@ while [[ $# -gt 0 ]]; do ORDERS="$2" shift 2 ;; + --ledger) + LEDGER="$2" + shift 2 + ;; + --chat-messages) + CHAT_MESSAGES="$2" + shift 2 + ;; -c|--concurrency) CONCURRENCY="$2" shift 2 @@ -115,6 +145,22 @@ while [[ $# -gt 0 ]]; do BASE_URL="$2" shift 2 ;; + --use-optimized) + USE_OPTIMIZED=true + shift + ;; + --use-improved) + USE_IMPROVED=true + shift + ;; + --gradual) + GRADUAL=true + shift + ;; + --warmup) + WARMUP="$2" + shift 2 + ;; -h|--help) show_usage exit 0 @@ -162,15 +208,45 @@ function check_backend() { # 生成测试数据 function generate_data() { log_info "开始生成测试数据..." - log_info "配置: 用户=$USERS, 商品=$LISTINGS, 订单=$ORDERS" + log_info "配置: 用户=$USERS, 商品=$LISTINGS, 订单=$ORDERS, 钱包流水=$LEDGER, 聊天消息=$CHAT_MESSAGES" check_database || exit 1 + # 选择使用的SQL脚本 + SQL_FILE="$SCRIPT_DIR/load_test_data.sql" + if [ "$USE_OPTIMIZED" = true ]; then + SQL_FILE="$SCRIPT_DIR/load_test_data_optimized.sql" + log_info "使用优化的数据生成脚本" + fi + # 创建临时SQL文件 TMP_SQL="/tmp/load_test_data_$(date +%s).sql" + trap 'rm -f "$TMP_SQL"' RETURN - cat > "$TMP_SQL" << EOF --- 临时生成的测试数据脚本 + if [ "$USE_OPTIMIZED" = true ]; then + cat > "$TMP_SQL" << EOF +-- 优化版数据生成 +USE hfb_sys; + +-- 调用批量生成存储过程 +CALL generate_users_batch($USERS); +CALL generate_listings_batch($LISTINGS); +CALL generate_orders_batch($ORDERS); +CALL generate_wallet_ledger_batch($LEDGER); +CALL generate_chat_data_batch($CHAT_MESSAGES); + +-- 显示统计 +SELECT '用户数' as item, COUNT(*) as count FROM users +UNION ALL SELECT '游戏账号数', COUNT(*) FROM game_accounts +UNION ALL SELECT '商品数', COUNT(*) FROM rental_listings +UNION ALL SELECT '订单数', COUNT(*) FROM rental_orders +UNION ALL SELECT '钱包流水数', COUNT(*) FROM wallet_ledger +UNION ALL SELECT '聊天会话数', COUNT(*) FROM chat_conversations +UNION ALL SELECT '聊天消息数', COUNT(*) FROM chat_messages; +EOF + else + cat > "$TMP_SQL" << EOF +-- 原始版数据生成 USE hfb_sys; -- 调用存储过程生成数据 @@ -186,7 +262,7 @@ FROM rental_orders WHERE id <= 1000 ON DUPLICATE KEY UPDATE order_id=order_id; -CALL generate_chat_messages(50000); +CALL generate_chat_messages($CHAT_MESSAGES); -- 显示统计 SELECT '用户数' as item, COUNT(*) as count FROM users @@ -197,16 +273,18 @@ UNION ALL SELECT '钱包流水数', COUNT(*) FROM wallet_ledger UNION ALL SELECT '聊天会话数', COUNT(*) FROM chat_conversations UNION ALL SELECT '聊天消息数', COUNT(*) FROM chat_messages; EOF + fi log_info "执行数据生成..." # 先执行基础SQL创建存储过程 - docker exec -i hfb-mysql mysql -uhfb -psecret hfb_sys < "$SCRIPT_DIR/load_test_data.sql" + docker exec -i hfb-mysql mysql -uhfb -psecret hfb_sys < "$SQL_FILE" # 执行数据生成 docker exec -i hfb-mysql mysql -uhfb -psecret hfb_sys < "$TMP_SQL" rm -f "$TMP_SQL" + trap - RETURN log_info "测试数据生成完成!" } @@ -221,11 +299,22 @@ function run_stress_test() { # 编译压测工具 log_info "编译压测工具..." cd "$SCRIPT_DIR" + go build -o stress_test stress_test.go + if [ $? -ne 0 ]; then + log_error "压测工具编译失败" + return 1 + fi + # 执行压测 log_info "开始执行压力测试..." - ./stress_test -url "$BASE_URL" -c "$CONCURRENCY" -d "$DURATION" -s "$SCENARIO" + + if [ "$GRADUAL" = true ]; then + ./stress_test -url "$BASE_URL" -c "$CONCURRENCY" -d "$DURATION" -s "$SCENARIO" -warmup "$WARMUP" -gradual + else + ./stress_test -url "$BASE_URL" -c "$CONCURRENCY" -d "$DURATION" -s "$SCENARIO" -warmup "$WARMUP" + fi log_info "压力测试完成!" }