更新每日推送数据的保存

This commit is contained in:
2026-07-22 11:00:42 +08:00
parent 140a068b88
commit b2b4fd4fbc
7 changed files with 621 additions and 61 deletions

View File

@@ -0,0 +1,481 @@
<?php
declare(strict_types=1);
namespace app\common\library;
use app\common\model\MallDailyPush;
use app\common\model\MallUserAsset;
use ba\Random;
use RuntimeException;
use support\think\Db;
use Throwable;
/**
* 将 daily_push_*.log 中保存的第三方原始推送幂等回放到业务表。
*/
class MallDailyPushLogReplay
{
private const QUERY_CHUNK_SIZE = 1000;
/**
* @return array<string, mixed>
*/
public function replay(string $logFile, bool $dryRun = false): array
{
$stats = $this->newStats($logFile, $dryRun);
$records = $this->parseRecords($logFile, $stats);
$stats['unique_members'] = count($records);
[$dailyMap, $assetByPlayxId, $assetByUsername] = $this->preloadExistingData($records);
$ratios = MallPlayxRatios::get();
$returnRatio = floatval($ratios['return_ratio'] ?? 0);
$unlockRatio = floatval($ratios['unlock_ratio'] ?? 0);
[$backupHandle, $backupPath] = $this->openBackup($dryRun);
$stats['backup_file'] = $backupPath;
try {
foreach ($records as $key => $record) {
try {
$daily = $dailyMap[$key] ?? null;
$asset = $assetByPlayxId[$record['user_id']]
?? ($record['username'] !== '' ? ($assetByUsername[$record['username']] ?? null) : null);
$result = $this->replayRecord(
$record,
$daily,
$asset,
$returnRatio,
$unlockRatio,
$dryRun,
$backupHandle,
$stats
);
if ($result['daily']) {
$dailyMap[$key] = $result['daily'];
}
if ($result['asset']) {
$assetByPlayxId[$record['user_id']] = $result['asset'];
if ($record['username'] !== '') {
$assetByUsername[$record['username']] = $result['asset'];
}
}
} catch (Throwable $e) {
$stats['failed_members']++;
$this->appendError($stats, 'user_id=' . $record['user_id'] . '' . $e->getMessage());
}
}
} finally {
if (is_resource($backupHandle)) {
fclose($backupHandle);
}
}
return $stats;
}
/**
* @return array<string, mixed>
*/
private function newStats(string $logFile, bool $dryRun): array
{
return [
'dry_run' => $dryRun,
'log_file' => $logFile,
'backup_file' => '',
'lines' => 0,
'members' => 0,
'unique_members' => 0,
'daily_created' => 0,
'daily_updated' => 0,
'daily_unchanged' => 0,
'assets_created' => 0,
'assets_updated' => 0,
'locked_points_delta' => 0,
'invalid_lines' => 0,
'failed_members' => 0,
'errors' => [],
];
}
/**
* @param array<string, mixed> $stats
* @return array<string, array<string, mixed>>
*/
private function parseRecords(string $logFile, array &$stats): array
{
if (!is_file($logFile) || !is_readable($logFile)) {
throw new RuntimeException('日志文件不存在或不可读:' . $logFile);
}
$handle = fopen($logFile, 'rb');
if ($handle === false) {
throw new RuntimeException('无法打开日志文件:' . $logFile);
}
$records = [];
try {
while (($line = fgets($handle)) !== false) {
$line = trim($line);
if ($line === '') {
continue;
}
$stats['lines']++;
$entry = json_decode($line, true);
$raw = is_array($entry) ? strval($entry['raw'] ?? '') : '';
$payload = $raw !== '' ? json_decode($raw, true) : null;
if (!is_array($payload) || !isset($payload['report_date'], $payload['member']) || !is_array($payload['member'])) {
$stats['invalid_lines']++;
$this->appendError($stats, '第 ' . strval($stats['lines']) . ' 行格式无效');
continue;
}
$date = $this->resolveDate($payload['report_date']);
if ($date === '') {
$stats['invalid_lines']++;
$this->appendError($stats, '第 ' . strval($stats['lines']) . ' 行 report_date 无效');
continue;
}
$entryTime = strtotime(strval($entry['time'] ?? ''));
$createTime = $entryTime === false ? time() : $entryTime;
foreach ($payload['member'] as $member) {
$stats['members']++;
if (!is_array($member)) {
$stats['failed_members']++;
continue;
}
$userId = trim(strval($member['member_id'] ?? ''));
if ($userId === '') {
$stats['failed_members']++;
$this->appendError($stats, '第 ' . strval($stats['lines']) . ' 行存在空 member_id');
continue;
}
$key = $date . '|' . $userId;
$records[$key] = [
'user_id' => $userId,
'date' => $date,
'username' => trim(strval($member['login'] ?? '')),
'phone' => trim(strval($member['phone'] ?? ($member['mobile'] ?? ''))),
'yesterday_win_loss_net' => $this->numericValue($member, ['yesterday_total_wl', 'yesterday_total_w']),
'yesterday_total_deposit' => $this->numericValue($member, ['yesterday_total_deposit']),
'lifetime_total_deposit' => $this->numericValue($member, ['ltv_deposit', 'lty_deposit']),
'lifetime_total_withdraw' => $this->numericValue($member, ['ltv_withdrawal', 'lty_withdrawal']),
'create_time' => $createTime,
];
}
}
} finally {
fclose($handle);
}
return $records;
}
/**
* @param array<string, array<string, mixed>> $records
* @return array{0:array<string,MallDailyPush>,1:array<string,MallUserAsset>,2:array<string,MallUserAsset>}
*/
private function preloadExistingData(array $records): array
{
$idsByDate = [];
$userIds = [];
$usernames = [];
foreach ($records as $record) {
$date = $record['date'];
$userId = $record['user_id'];
$idsByDate[$date][$userId] = true;
$userIds[$userId] = true;
if ($record['username'] !== '') {
$usernames[$record['username']] = true;
}
}
$dailyMap = [];
foreach ($idsByDate as $date => $ids) {
foreach (array_chunk(array_keys($ids), self::QUERY_CHUNK_SIZE) as $idChunk) {
$rows = MallDailyPush::where('date', $date)->whereIn('user_id', $idChunk)->select();
foreach ($rows as $row) {
/** @var MallDailyPush $row */
$dailyMap[$date . '|' . strval($row->user_id)] = $row;
}
}
}
$assetByPlayxId = [];
foreach (array_chunk(array_keys($userIds), self::QUERY_CHUNK_SIZE) as $idChunk) {
$rows = MallUserAsset::whereIn('playx_user_id', $idChunk)->select();
foreach ($rows as $row) {
$assetByPlayxId[strval($row->playx_user_id)] = $row;
}
}
$assetByUsername = [];
foreach (array_chunk(array_keys($usernames), self::QUERY_CHUNK_SIZE) as $usernameChunk) {
$rows = MallUserAsset::whereIn('username', $usernameChunk)->select();
foreach ($rows as $row) {
$assetByUsername[strval($row->username)] = $row;
}
}
return [$dailyMap, $assetByPlayxId, $assetByUsername];
}
/**
* @return array{0:mixed,1:string}
*/
private function openBackup(bool $dryRun): array
{
if ($dryRun) {
return [null, ''];
}
$backupDir = runtime_path('backup/daily_push_replay');
if (!is_dir($backupDir) && !mkdir($backupDir, 0755, true) && !is_dir($backupDir)) {
throw new RuntimeException('无法创建回放备份目录');
}
$path = $backupDir . DIRECTORY_SEPARATOR . 'before_replay_' . date('Ymd_His') . '.jsonl';
$handle = fopen($path, 'wb');
if ($handle === false) {
throw new RuntimeException('无法创建回放备份文件');
}
return [$handle, $path];
}
/**
* @param array<string, mixed> $record
* @param resource|null $backupHandle
* @param array<string, mixed> $stats
* @return array{daily:?MallDailyPush,asset:?MallUserAsset}
*/
private function replayRecord(
array $record,
?MallDailyPush $daily,
?MallUserAsset $asset,
float $returnRatio,
float $unlockRatio,
bool $dryRun,
$backupHandle,
array &$stats
): array {
$oldWinLossNet = $daily ? floatval($daily->yesterday_win_loss_net ?? 0) : 0.0;
$lockedDelta = $this->lockedContribution($record['yesterday_win_loss_net'], $returnRatio)
- $this->lockedContribution($oldWinLossNet, $returnRatio);
$todayLimit = intval(round($record['yesterday_total_deposit'] * $unlockRatio));
$dailyData = [
'user_id' => $record['user_id'],
'date' => $record['date'],
'username' => $record['username'],
'yesterday_win_loss_net' => $record['yesterday_win_loss_net'],
'yesterday_total_deposit' => $record['yesterday_total_deposit'],
'lifetime_total_deposit' => $record['lifetime_total_deposit'],
'lifetime_total_withdraw' => $record['lifetime_total_withdraw'],
];
$dailyChanged = !$daily || $this->dailyHasChanges($daily, $dailyData);
$this->countDailyChange($stats, $daily, $dailyChanged);
$assetWillCreate = !$asset;
$assetWillUpdate = $assetWillCreate
|| trim(strval($asset->playx_user_id ?? '')) !== $record['user_id']
|| $lockedDelta !== 0
|| ($record['username'] !== '' && trim(strval($asset->username ?? '')) !== $record['username'])
|| ($record['phone'] !== '' && trim(strval($asset->phone ?? '')) === '')
|| $this->shouldUpdateDailyLimit($asset, $record['date'], $todayLimit);
if ($assetWillCreate) {
$stats['assets_created']++;
} elseif ($assetWillUpdate) {
$stats['assets_updated']++;
}
$stats['locked_points_delta'] += $lockedDelta;
if ($dryRun) {
return ['daily' => $daily, 'asset' => $asset];
}
if (!$dailyChanged && !$assetWillUpdate) {
return ['daily' => $daily, 'asset' => $asset];
}
$this->writeBackup($backupHandle, $record['date'] . '|' . $record['user_id'], $daily, $asset);
Db::startTrans();
try {
if ($daily) {
if ($dailyChanged) {
$daily->save($dailyData);
}
} else {
$dailyData['create_time'] = $record['create_time'];
$daily = MallDailyPush::create($dailyData);
}
if (!$asset) {
$effectiveUsername = $record['username'] !== '' ? $record['username'] : 'playx_' . $record['user_id'];
$asset = MallUserAsset::create([
'playx_user_id' => $record['user_id'],
'username' => $effectiveUsername,
'phone' => $record['phone'],
'password' => hash_password(Random::build('alnum', 16)),
'admin_id' => 0,
'locked_points' => 0,
'available_points' => 0,
'today_limit' => 0,
'today_claimed' => 0,
'today_limit_date' => null,
'create_time' => $record['create_time'],
'update_time' => time(),
]);
if (!$asset) {
throw new RuntimeException('创建用户资产失败');
}
/** @var MallUserAsset $asset */
} elseif ($assetWillUpdate) {
$asset->playx_user_id = $record['user_id'];
if ($record['username'] !== '') {
$asset->username = $record['username'];
}
if ($record['phone'] !== '' && trim(strval($asset->phone ?? '')) === '') {
$asset->phone = $record['phone'];
}
}
if ($assetWillUpdate) {
if ($lockedDelta !== 0) {
$asset->locked_points = max(0, intval($asset->locked_points ?? 0) + $lockedDelta);
}
$this->applyDailyLimit($asset, $record['date'], $todayLimit);
$asset->save();
}
Db::commit();
} catch (Throwable $e) {
Db::rollback();
throw $e;
}
return ['daily' => $daily, 'asset' => $asset];
}
/**
* @param array<string, mixed> $stats
*/
private function countDailyChange(array &$stats, ?MallDailyPush $daily, bool $changed): void
{
if (!$daily) {
$stats['daily_created']++;
} elseif ($changed) {
$stats['daily_updated']++;
} else {
$stats['daily_unchanged']++;
}
}
/**
* @param array<string, mixed> $member
* @param array<int, string> $keys
*/
private function numericValue(array $member, array $keys): float
{
foreach ($keys as $key) {
if (array_key_exists($key, $member) && is_numeric($member[$key])) {
return round(floatval($member[$key]), 2);
}
}
return 0.0;
}
private function resolveDate(mixed $reportDate): string
{
if (is_numeric($reportDate)) {
$timestamp = intval($reportDate);
return $timestamp > 0 ? date('Y-m-d', $timestamp) : '';
}
$date = trim(strval($reportDate));
$parsed = \DateTime::createFromFormat('Y-m-d', $date);
return $parsed && $parsed->format('Y-m-d') === $date ? $date : '';
}
private function lockedContribution(float $winLossNet, float $returnRatio): int
{
return $winLossNet < 0 ? intval(round(abs($winLossNet) * $returnRatio)) : 0;
}
/**
* @param array<string, mixed> $dailyData
*/
private function dailyHasChanges(MallDailyPush $daily, array $dailyData): bool
{
if (trim(strval($daily->username ?? '')) !== $dailyData['username']) {
return true;
}
foreach ([
'yesterday_win_loss_net',
'yesterday_total_deposit',
'lifetime_total_deposit',
'lifetime_total_withdraw',
] as $field) {
if (abs(floatval($daily->$field ?? 0) - floatval($dailyData[$field])) > 0.000001) {
return true;
}
}
return false;
}
private function shouldUpdateDailyLimit(?MallUserAsset $asset, string $date, int $todayLimit): bool
{
if (!$asset) {
return true;
}
$currentDate = trim(strval($asset->today_limit_date ?? ''));
if ($currentDate !== '' && $currentDate > $date) {
return false;
}
return $currentDate !== $date || intval($asset->today_limit ?? 0) !== $todayLimit;
}
private function applyDailyLimit(MallUserAsset $asset, string $date, int $todayLimit): void
{
$currentDate = trim(strval($asset->today_limit_date ?? ''));
if ($currentDate !== '' && $currentDate > $date) {
return;
}
if ($currentDate !== $date) {
$asset->today_claimed = 0;
$asset->today_limit_date = $date;
}
$asset->today_limit = $todayLimit;
}
/**
* @param resource|null $handle
*/
private function writeBackup($handle, string $key, ?MallDailyPush $daily, ?MallUserAsset $asset): void
{
if (!is_resource($handle)) {
return;
}
$assetData = null;
if ($asset) {
$assetData = [
'id' => $asset->id,
'playx_user_id' => $asset->playx_user_id,
'username' => $asset->username,
'phone' => $asset->phone,
'locked_points' => $asset->locked_points,
'today_limit' => $asset->today_limit,
'today_claimed' => $asset->today_claimed,
'today_limit_date' => $asset->today_limit_date,
'update_time' => $asset->update_time,
];
}
fwrite($handle, json_encode([
'key' => $key,
'daily_push' => $daily ? $daily->toArray() : null,
'user_asset' => $assetData,
], JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES) . PHP_EOL);
}
/**
* @param array<string, mixed> $stats
*/
private function appendError(array &$stats, string $message): void
{
if (count($stats['errors']) < 100) {
$stats['errors'][] = $message;
}
}
}