fix: report recharge in fen with daily three-hour cutoff
This commit is contained in:
parent
2224688b65
commit
d137b4e030
|
|
@ -2,12 +2,14 @@ import cloud from '@lafjs/cloud';
|
|||
|
||||
const DAY = 24 * 60 * 60 * 1000;
|
||||
const LEGACY_TIME_OFFSET = 8 * 60 * 60 * 1000;
|
||||
const REPORT_DELAY = 3 * 60 * 60 * 1000;
|
||||
const BATCH_SIZE = 100;
|
||||
|
||||
function emptyStats(asOf: number) {
|
||||
return {
|
||||
version: 1,
|
||||
version: 2,
|
||||
currency: 'CNY',
|
||||
unit: 'fen',
|
||||
asOf,
|
||||
amount15d: 0,
|
||||
amount30d: 0,
|
||||
|
|
@ -42,7 +44,8 @@ function timestamp(value: any): number | null {
|
|||
export default async function () {
|
||||
const mongo = cloud.mongo.db;
|
||||
const users = mongo.collection('users');
|
||||
const asOf = Date.now();
|
||||
// The 03:00 run reports through 00:00; this is separate from the stored +8h correction.
|
||||
const asOf = Date.now() - REPORT_DELAY;
|
||||
let afterId: any;
|
||||
let processedUsers = 0;
|
||||
let updatedUsers = 0;
|
||||
|
|
@ -90,7 +93,7 @@ export default async function () {
|
|||
// order.time is the original, unshifted Unix timestamp.
|
||||
const chargeTime = timestamp(order.chargeTime);
|
||||
const paidAt = chargeTime === null ? timestamp(order.time) : chargeTime - LEGACY_TIME_OFFSET;
|
||||
if (paidAt !== null && paidAt > asOf) continue;
|
||||
if (paidAt !== null && paidAt >= asOf) continue;
|
||||
if (summary.seen.has(order.outTradeNo)) {
|
||||
summary.stats.duplicateOrderCount++;
|
||||
continue;
|
||||
|
|
@ -112,9 +115,9 @@ export default async function () {
|
|||
}
|
||||
|
||||
for (const summary of summaries.values()) {
|
||||
summary.stats.amount15d = summary.fen15d / 100;
|
||||
summary.stats.amount30d = summary.fen30d / 100;
|
||||
summary.stats.amountTotal = summary.fenTotal / 100;
|
||||
summary.stats.amount15d = summary.fen15d;
|
||||
summary.stats.amount30d = summary.fen30d;
|
||||
summary.stats.amountTotal = summary.fenTotal;
|
||||
}
|
||||
const result = await users.bulkWrite(batch.map(user => {
|
||||
const rechargeStats = summaries.get(user.openid)?.stats ?? { ...emptyStats(asOf), missingOpenid: true };
|
||||
|
|
@ -124,7 +127,12 @@ export default async function () {
|
|||
_id: user._id,
|
||||
openid: user.openid ?? null,
|
||||
pay_user: true,
|
||||
$or: [{ 'rechargeStats.asOf': { $exists: false } }, { 'rechargeStats.asOf': { $lte: asOf } }],
|
||||
$or: [
|
||||
// Version 1 used yuan and a later cutoff: replace it even when its asOf is newer.
|
||||
{ 'rechargeStats.version': 1 },
|
||||
{ 'rechargeStats.asOf': { $exists: false } },
|
||||
{ 'rechargeStats.asOf': { $lte: asOf } },
|
||||
],
|
||||
},
|
||||
update: { $set: { rechargeStats } },
|
||||
} };
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
name: rechargeStats
|
||||
desc: "付费用户近15天、近30天及累计充值统计;每小时由定时触发器执行"
|
||||
desc: "付费用户充值统计(分);每日03:00执行,截止时间回退3小时"
|
||||
methods: []
|
||||
tags:
|
||||
- statistics
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
{
|
||||
"desc": "付费用户充值统计(每小时)",
|
||||
"desc": "付费用户充值统计(每日凌晨3点)",
|
||||
"target": "rechargeStats",
|
||||
"cron": "0 * * * *"
|
||||
"cron": "0 3 * * *"
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,34 +1,38 @@
|
|||
# 付费用户充值统计定时任务
|
||||
|
||||
云函数 `rechargeStats` 每次全量重算当前 `users.pay_user === true` 的玩家,将结果覆盖到 `users.rechargeStats`。默认每小时整点执行;不修改订单和玩家付费标记。
|
||||
云函数 `rechargeStats` 每次全量重算当前 `users.pay_user === true` 的玩家,将结果覆盖到 `users.rechargeStats`。每日凌晨 03:00 执行,统计截止时间固定回退 3 小时;不修改订单和玩家付费标记。
|
||||
|
||||
## 保存字段
|
||||
|
||||
| 字段 | 含义 |
|
||||
| --- | --- |
|
||||
| `amount15d` | 从本轮运行时刻倒推 15 × 24 小时的充值金额,元 |
|
||||
| `amount30d` | 从本轮运行时刻倒推 30 × 24 小时的充值金额,元 |
|
||||
| `amountTotal` | 截至本轮运行时刻、当前库中可统计的累计充值金额,元 |
|
||||
| `amount15d` | 从统计截止时间倒推 15 × 24 小时的充值金额,整数分 |
|
||||
| `amount30d` | 从统计截止时间倒推 30 × 24 小时的充值金额,整数分 |
|
||||
| `amountTotal` | 统计截止时间之前、当前库中可统计的累计充值金额,整数分 |
|
||||
| `currency` | 固定 `CNY` |
|
||||
| `asOf` | 本轮统计截止时间,原始 Unix 毫秒时间戳 |
|
||||
| `unit` | 固定 `fen`,金额单位为分 |
|
||||
| `asOf` | 本轮运行时间减 3 小时,原始 Unix 毫秒时间戳 |
|
||||
| `orderCount` | 纳入累计金额的去重订单数 |
|
||||
| `fallbackTimeOrderCount` | 无支付确认时间,使用下单时间的订单数 |
|
||||
| `missingTimeOrderCount` | 支付和下单时间均缺失,仅计入累计的订单数 |
|
||||
| `invalidOrderCount` | 缺失订单号、金额或数量无效而跳过的记录数 |
|
||||
| `duplicateOrderCount` | 跳过的重复订单记录数 |
|
||||
| `missingOpenid` | 用户是否缺少可关联订单的 openid |
|
||||
| `version` | 统计规则版本,当前为 1 |
|
||||
| `version` | 统计规则版本,当前为 2 |
|
||||
|
||||
版本 1 的金额单位为元。版本 2 首次运行会重新计算并整体覆盖旧统计,即使旧版 `asOf` 晚于新版截止时间,也允许完成升级。读取方应以 `version: 2 / unit: "fen"` 识别分单位;运行未覆盖到的用户仍可能保留旧版数据,不能仅按字段名判断单位。
|
||||
|
||||
## 统计口径
|
||||
|
||||
- 通过 `users.openid = order.openid` 关联,仅纳入数值型 `state: 1/2` 的记录。未确认支付的 `state: 0` 和临时集合 `iosOrder` 不参与。
|
||||
- 排除 `paymentAppEnv: "test"` 或 `outTradeNo` 以 `wct_` 开头的测试订单;保留没有环境字段的历史订单。
|
||||
- 每个 openid 内按 `outTradeNo` 去重,同号多条按 `_id` 升序取第一条金额有效、时间不在未来的记录。若同号金额不同,应人工核查重复数据。
|
||||
- 金额为 `goodsPrice × itemCount`,先用整数分累加,再除以 100 保存为元。金额和数量兼容数字字符串,要求正安全整数;缺失数量不默认当作 1,异常记录计入 `invalidOrderCount`。累计金额超出安全整数范围则报错,不保存失真的数值。
|
||||
- 每轮只获取一次当前时间,窗口为 `[asOf - N × 86400000, asOf]`,不是自然日,也不是从午夜倒推。
|
||||
- 每个 openid 内按 `outTradeNo` 去重,同号多条按 `_id` 升序取第一条金额有效、未达到统计截止时间的记录。若同号金额不同,应人工核查重复数据。
|
||||
- 金额为 `goodsPrice × itemCount`,用整数分累加并直接保存,不除以 100。金额和数量兼容数字字符串,要求正安全整数;缺失数量不默认当作 1,异常记录计入 `invalidOrderCount`。累计金额超出安全整数范围则报错,不保存失真的数值。
|
||||
- 每轮只获取一次当前时间,`asOf = Date.now() - 3小时`,窗口为 `[asOf - N × 86400000, asOf)`,含起点、不含截止点。累计金额也使用同一截止点。每天 03:00 准时运行时,截止点即当天 00:00,凌晨 00:00~03:00 的订单留到下一天统计。手动运行或触发延迟时也固定回退 3 小时,不自动截断到午夜。
|
||||
- 例如北京时间 2026-09-15 03:00 执行:15 天窗口为 `[2026-08-31 00:00, 2026-09-15 00:00)`;30 天窗口为 `[2026-08-16 00:00, 2026-09-15 00:00)`;累计统计到 2026-09-15 00:00 之前。
|
||||
- 优先使用 `chargeTime`。当前支付写入代码使用 `new Date(Date.now() + 8小时)`,因此还原时减去 8 小时;此处针对现有存储约定,若以后修正支付时间存储方式,必须同步调整统计规则。
|
||||
- `chargeTime` 缺失或无效时回退到原始 `order.time`,不再减 8 小时;两个时间都不可用时只计入累计,并记录异常数量。支持原始毫秒数、Date 和带时区的 ISO 字符串;不解析无时区日期字符串。
|
||||
- 已知时间在未来的订单暂不纳入任何金额。没有匹配订单的付费用户也保存零值结果;这不代表迁移前从未充值,结合 `orderCount` 核查历史完整性。
|
||||
- 已知时间达到或晚于 `asOf` 的订单暂不纳入任何金额。没有匹配订单的付费用户也保存零值结果;这不代表迁移前从未充值,结合 `orderCount` 核查历史完整性。
|
||||
- 当前项目没有退款同步逻辑,因此这是当前已确认订单的金额统计,不扣除退款,不还原渠道优惠后的实付。订单状态异常仍可能造成偏差。
|
||||
|
||||
## 执行与一致性
|
||||
|
|
@ -50,12 +54,12 @@ db.collection('order').createIndex({ openid: 1, outTradeNo: 1, _id: 1 });
|
|||
|
||||
1. 在目标 Laf 应用发布 `functions/rechargeStats.ts` 和同名 YAML,保持 `methods: []`,不开放公共 HTTP 入口。
|
||||
2. 在云函数控制台手动执行一次 `rechargeStats`(无参数),检查返回的 `processedUsers`、`updatedUsers` 和部分玩家的结果,确认运行耗时。
|
||||
3. 在触发器面板绑定函数 `rechargeStats`,使用 `recharge-stats.trigger.json` 的配置:每小时整点执行。
|
||||
3. 在触发器面板绑定函数 `rechargeStats`,使用 `recharge-stats.trigger.json` 的配置:每日 03:00 执行。表达式 `0 3 * * *` 按触发器时区解释;目标是北京时间 03:00,应确认调度时区为 `Asia/Shanghai`。如果部署使用 UTC 调度,则使用 `0 19 * * *`(UTC 19:00 为次日北京时间 03:00)。3 小时统计回退和订单存储的 8 小时修正是两件事,不改变 Unix 时间戳所属时区。
|
||||
4. 检查已有触发器,避免重复创建。若使用已登录且已绑定正确应用的 Laf CLI,可执行:
|
||||
|
||||
```shell
|
||||
laf trigger list
|
||||
laf trigger create "付费用户充值统计(每小时)" rechargeStats "0 * * * *"
|
||||
laf trigger create "付费用户充值统计(每日凌晨3点)" rechargeStats "0 3 * * *"
|
||||
```
|
||||
|
||||
5. 首次定时触发后确认日志出现 `rechargeStats completed`,抽查 `users.rechargeStats.asOf` 已更新。大规模历史数据上线前应确认全量耗时在云函数执行时限内。
|
||||
|
|
@ -70,4 +74,4 @@ Node.js 24.11 或兼容 `registerHooks` / `stripTypeScriptTypes` 的版本:
|
|||
node --test laf-cloud/tests/recharge-stats.test.mjs
|
||||
```
|
||||
|
||||
覆盖滚动时间边界、8 小时修正、支付时间优先、时间回退、无效金额、数量、去重、测试订单排除、超过 1000 条订单、多页用户、重跑、并发写保护、异常中断和关闭 HTTP 入口。使用内存 MongoDB 接口替身,不连接生产数据库。
|
||||
覆盖分单位、3 小时截止回退、午夜半开区间边界、凌晨订单次日计入、旧版元统计升级、8 小时存储修正、支付时间优先、时间回退、无效金额、数量、去重、测试订单排除、超过 1000 条订单、多页用户、重跑、并发写保护、异常中断和关闭 HTTP 入口。使用内存 MongoDB 接口替身,不连接生产数据库。
|
||||
|
|
|
|||
|
|
@ -5,7 +5,8 @@ import { readFileSync } from 'node:fs';
|
|||
|
||||
const DAY = 86400000;
|
||||
const OFFSET = 8 * 3600000;
|
||||
const NOW = Date.UTC(2026, 8, 15, 12);
|
||||
const NOW = Date.parse('2026-09-15T03:00:00+08:00');
|
||||
const CUTOFF = NOW - 3 * 3600000;
|
||||
const state = { users: [], order: [], writes: 0, closed: 0, beforeWrite: null, failRead: false };
|
||||
const realNow = Date.now;
|
||||
Date.now = () => NOW;
|
||||
|
|
@ -93,23 +94,26 @@ function reset() {
|
|||
function order(id, age, extra = {}) {
|
||||
return {
|
||||
_id: id, outTradeNo: id, openid: 'o1', state: 2, goodsPrice: 100, itemCount: 1,
|
||||
chargeTime: new Date(NOW - age + OFFSET), time: NOW - age - DAY, ...extra,
|
||||
chargeTime: new Date(CUTOFF - age + OFFSET), time: CUTOFF - age - DAY, ...extra,
|
||||
};
|
||||
}
|
||||
|
||||
test('rolling windows include exact boundaries, use payment time, restore +8h and exclude future payments', async () => {
|
||||
test('03:00 run uses midnight cutoff for all totals, including the start and excluding the end', async () => {
|
||||
reset();
|
||||
state.order = [
|
||||
order('now', 0), order('15d', 15 * DAY), order('15d-old', 15 * DAY + 1),
|
||||
order('30d', 30 * DAY), order('30d-old', 30 * DAY + 1), order('future', -1),
|
||||
order('before-midnight', 1), order('15d', 15 * DAY), order('15d-old', 15 * DAY + 1),
|
||||
order('30d', 30 * DAY), order('30d-old', 30 * DAY + 1), order('midnight', 0),
|
||||
order('early-morning', -2 * 3600000), order('future', -4 * 3600000),
|
||||
];
|
||||
await run();
|
||||
const stats = state.users[0].rechargeStats;
|
||||
assert.equal(stats.amount15d, 2);
|
||||
assert.equal(stats.amount30d, 4);
|
||||
assert.equal(stats.amountTotal, 5);
|
||||
assert.equal(stats.amount15d, 200);
|
||||
assert.equal(stats.amount30d, 400);
|
||||
assert.equal(stats.amountTotal, 500);
|
||||
assert.equal(stats.orderCount, 5);
|
||||
assert.equal(stats.asOf, NOW);
|
||||
assert.equal(stats.asOf, Date.parse('2026-09-15T00:00:00+08:00'));
|
||||
assert.equal(stats.version, 2);
|
||||
assert.equal(stats.unit, 'fen');
|
||||
});
|
||||
|
||||
test('counts quantities, numeric strings, paid states and deduplicates; excludes unpaid and test orders', async () => {
|
||||
|
|
@ -122,7 +126,7 @@ test('counts quantities, numeric strings, paid states and deduplicates; excludes
|
|||
order('test-env', DAY, { paymentAppEnv: 'test' }), order('wct_test', DAY),
|
||||
];
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 3.10);
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 310);
|
||||
assert.equal(state.users[0].rechargeStats.orderCount, 2);
|
||||
assert.equal(state.users[0].rechargeStats.duplicateOrderCount, 1);
|
||||
});
|
||||
|
|
@ -130,18 +134,18 @@ test('counts quantities, numeric strings, paid states and deduplicates; excludes
|
|||
test('falls back to unshifted creation time, tracks unknown times, and rejects malformed amounts', async () => {
|
||||
reset();
|
||||
state.order = [
|
||||
order('fallback', 0, { chargeTime: 0, time: NOW - 15 * DAY }),
|
||||
order('fallback', 0, { chargeTime: 0, time: CUTOFF - 15 * DAY }),
|
||||
order('unknown', 0, { chargeTime: 0, time: null }),
|
||||
order('iso', 0, { chargeTime: new Date(NOW - 30 * DAY + OFFSET).toISOString() }),
|
||||
order('iso', 0, { chargeTime: new Date(CUTOFF - 30 * DAY + OFFSET).toISOString() }),
|
||||
order('bad-price', 0, { goodsPrice: 'bad' }), order('bad-count', 0, { itemCount: null }),
|
||||
order('negative', 0, { goodsPrice: -1 }), order('missing-id', 0, { outTradeNo: '' }),
|
||||
order('overflow', 0, { goodsPrice: Number.MAX_SAFE_INTEGER, itemCount: 2 }),
|
||||
];
|
||||
await run();
|
||||
const stats = state.users[0].rechargeStats;
|
||||
assert.equal(stats.amount15d, 1);
|
||||
assert.equal(stats.amount30d, 2);
|
||||
assert.equal(stats.amountTotal, 3);
|
||||
assert.equal(stats.amount15d, 100);
|
||||
assert.equal(stats.amount30d, 200);
|
||||
assert.equal(stats.amountTotal, 300);
|
||||
assert.equal(stats.fallbackTimeOrderCount, 1);
|
||||
assert.equal(stats.missingTimeOrderCount, 1);
|
||||
assert.equal(stats.invalidOrderCount, 5);
|
||||
|
|
@ -156,11 +160,11 @@ test('processes every page, restricts to boolean true, resets old totals, and fl
|
|||
state.users.push({ _id: 'no-openid', pay_user: true });
|
||||
state.users.push({ _id: 'free', openid: 'free', pay_user: false });
|
||||
state.users.push({ _id: 'string', openid: 'string', pay_user: 'true' });
|
||||
state.order = [order('last', 0, { openid: 'o250' }), order('free', 0, { openid: 'free' })];
|
||||
state.order = [order('last', 1, { openid: 'o250' }), order('free', 1, { openid: 'free' })];
|
||||
const result = await run();
|
||||
assert.equal(result.data.processedUsers, 252);
|
||||
assert.equal(state.writes, 3);
|
||||
assert.equal(state.users.find(user => user.openid === 'o250').rechargeStats.amountTotal, 1);
|
||||
assert.equal(state.users.find(user => user.openid === 'o250').rechargeStats.amountTotal, 100);
|
||||
assert.equal(state.users.find(user => user.openid === 'o0').rechargeStats.amountTotal, 0);
|
||||
assert.equal(state.users.find(user => user._id === 'no-openid').rechargeStats.missingOpenid, true);
|
||||
assert.equal(state.users.find(user => user._id === 'free').rechargeStats, undefined);
|
||||
|
|
@ -169,26 +173,26 @@ test('processes every page, restricts to boolean true, resets old totals, and fl
|
|||
|
||||
test('reruns overwrite instead of accumulating; older runs cannot overwrite newer snapshots', async () => {
|
||||
reset();
|
||||
state.order = [order('paid', 0)];
|
||||
state.order = [order('paid', 1)];
|
||||
await run();
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 1);
|
||||
state.users[0].rechargeStats = { asOf: NOW + 1, amountTotal: 2 };
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 100);
|
||||
state.users[0].rechargeStats = { version: 2, asOf: CUTOFF + 1, amountTotal: 200 };
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 2);
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 200);
|
||||
});
|
||||
|
||||
test('streams more than 1000 orders and expired windows clear on the next run', async () => {
|
||||
reset();
|
||||
state.order = Array.from({ length: 1005 }, (_, i) => order(`paid-${i}`, 30 * DAY, { goodsPrice: 1 }));
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 10.05);
|
||||
assert.equal(state.users[0].rechargeStats.amount30d, 10.05);
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 1005);
|
||||
assert.equal(state.users[0].rechargeStats.amount30d, 1005);
|
||||
try {
|
||||
Date.now = () => NOW + 1;
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amount30d, 0);
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 10.05);
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 1005);
|
||||
} finally { Date.now = () => NOW; }
|
||||
});
|
||||
|
||||
|
|
@ -218,5 +222,33 @@ test('timer function does not expose a public HTTP endpoint', () => {
|
|||
assert.match(config, /methods: \[\]/);
|
||||
const trigger = JSON.parse(readFileSync(new URL('../recharge-stats.trigger.json', import.meta.url), 'utf8'));
|
||||
assert.equal(trigger.target, 'rechargeStats');
|
||||
assert.equal(trigger.cron, '0 * * * *');
|
||||
assert.equal(trigger.cron, '0 3 * * *');
|
||||
});
|
||||
|
||||
test('first version 2 run replaces version 1 yuan snapshot despite its later cutoff', async () => {
|
||||
reset();
|
||||
state.users[0].rechargeStats = { version: 1, asOf: NOW - 1, amount15d: 1, amount30d: 1, amountTotal: 1 };
|
||||
state.order = [order('paid', DAY, { goodsPrice: 101, itemCount: 3 })];
|
||||
await run();
|
||||
const stats = state.users[0].rechargeStats;
|
||||
assert.equal(stats.version, 2);
|
||||
assert.equal(stats.unit, 'fen');
|
||||
assert.equal(stats.asOf, CUTOFF);
|
||||
assert.equal(stats.amount15d, 303);
|
||||
assert.equal(stats.amount30d, 303);
|
||||
assert.equal(stats.amountTotal, 303);
|
||||
});
|
||||
|
||||
test('payments between midnight and 03:00 are included on the following daily run', async () => {
|
||||
reset();
|
||||
state.order = [order('midnight', 0), order('02:00', -2 * 3600000)];
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 0);
|
||||
try {
|
||||
Date.now = () => NOW + DAY;
|
||||
await run();
|
||||
assert.equal(state.users[0].rechargeStats.amount15d, 200);
|
||||
assert.equal(state.users[0].rechargeStats.amount30d, 200);
|
||||
assert.equal(state.users[0].rechargeStats.amountTotal, 200);
|
||||
} finally { Date.now = () => NOW; }
|
||||
});
|
||||
|
|
|
|||
Loading…
Reference in New Issue
Block a user