From b5d46b238161cf726b7ff859589cb17ae1b085be Mon Sep 17 00:00:00 2001 From: yml2213 Date: Thu, 20 Aug 2026 23:04:52 +0800 Subject: [PATCH] =?UTF-8?q?=E5=85=A8=E9=A1=B9=E7=9B=AE=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E5=BA=93=E6=80=A7=E8=83=BD=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.mac-docker.example | 4 + .env.server.example | 4 + apps/backend/src/config/defaults.ts | 5 +- apps/backend/src/config/env-overrides.ts | 3 + apps/backend/src/config/runtime-validation.ts | 29 +++++++ apps/backend/src/db/client.ts | 3 + .../048_database_performance_guards.sql | 13 ++++ .../src/repositories/realtime-event-repo.ts | 16 ++-- .../worker-platform/work-order-repo.ts | 76 +++++++++++++------ .../services/admin/admin-dashboard-service.ts | 48 +++++++++--- .../worker-platform/worker-service.ts | 2 - apps/backend/src/types/runtime-config.ts | 3 + docker-compose.yml | 5 +- 13 files changed, 168 insertions(+), 43 deletions(-) create mode 100644 apps/backend/src/db/migrations/048_database_performance_guards.sql diff --git a/.env.mac-docker.example b/.env.mac-docker.example index 8a942e1f..5f0022fd 100644 --- a/.env.mac-docker.example +++ b/.env.mac-docker.example @@ -16,6 +16,10 @@ POSTGRES_DB=order_site POSTGRES_USER=postgres POSTGRES_PASSWORD=postgres DATABASE_URL=postgres://postgres:postgres@postgres:5432/order_site +DATABASE_MAX_CONNECTIONS=8 +DATABASE_IDLE_TIMEOUT_MS=30000 +DATABASE_CONNECTION_TIMEOUT_MS=5000 +DATABASE_STATEMENT_TIMEOUT_MS=15000 POSTGRES_PORT=5432 # Logging diff --git a/.env.server.example b/.env.server.example index 21c62876..14039e98 100644 --- a/.env.server.example +++ b/.env.server.example @@ -23,6 +23,10 @@ POSTGRES_DB=order_site POSTGRES_USER=postgres POSTGRES_PASSWORD=change-me-postgres-password DATABASE_URL=postgres://postgres:change-me-postgres-password@postgres:5432/order_site +DATABASE_MAX_CONNECTIONS=8 +DATABASE_IDLE_TIMEOUT_MS=30000 +DATABASE_CONNECTION_TIMEOUT_MS=5000 +DATABASE_STATEMENT_TIMEOUT_MS=15000 # Logging LOG_LEVEL=info diff --git a/apps/backend/src/config/defaults.ts b/apps/backend/src/config/defaults.ts index 97a00ae6..0d7d9724 100644 --- a/apps/backend/src/config/defaults.ts +++ b/apps/backend/src/config/defaults.ts @@ -22,7 +22,10 @@ export function createDefaultRuntimeConfig(projectRoot: string): RuntimeConfig { database: { url: 'postgres://postgres:postgres@127.0.0.1:5432/order_site', ssl: false, - maxConnections: 10, + maxConnections: 8, + idleTimeoutMs: 30_000, + connectionTimeoutMs: 5_000, + statementTimeoutMs: 15_000, }, storage: { diff --git a/apps/backend/src/config/env-overrides.ts b/apps/backend/src/config/env-overrides.ts index 9dafb1e6..dc61532c 100644 --- a/apps/backend/src/config/env-overrides.ts +++ b/apps/backend/src/config/env-overrides.ts @@ -35,6 +35,9 @@ export const ENV_OVERRIDES: readonly EnvOverride[] = [ stringEnv('DATABASE_URL', ['database', 'url']), booleanEnv('DATABASE_SSL', ['database', 'ssl']), integerEnv('DATABASE_MAX_CONNECTIONS', ['database', 'maxConnections']), + integerEnv('DATABASE_IDLE_TIMEOUT_MS', ['database', 'idleTimeoutMs']), + integerEnv('DATABASE_CONNECTION_TIMEOUT_MS', ['database', 'connectionTimeoutMs']), + integerEnv('DATABASE_STATEMENT_TIMEOUT_MS', ['database', 'statementTimeoutMs']), stringEnv('STORAGE_ENDPOINT', ['storage', 'endpoint']), stringEnv('STORAGE_BUCKET', ['storage', 'bucket']), stringEnv('STORAGE_ACCESS_KEY_ID', ['storage', 'accessKeyId']), diff --git a/apps/backend/src/config/runtime-validation.ts b/apps/backend/src/config/runtime-validation.ts index 985ad346..43bdc468 100644 --- a/apps/backend/src/config/runtime-validation.ts +++ b/apps/backend/src/config/runtime-validation.ts @@ -45,6 +45,23 @@ export function validateRuntimeConfig( requireInteger(issues, 'server.port', config.server?.port, { min: 1, max: 65535 }) requireInteger(issues, 'database.maxConnections', config.database?.maxConnections, { min: 1 }) + requireOptionalInteger(issues, 'database.idleTimeoutMs', config.database?.idleTimeoutMs, { + min: 1, + }) + requireOptionalInteger( + issues, + 'database.connectionTimeoutMs', + config.database?.connectionTimeoutMs, + { min: 1 }, + ) + requireOptionalInteger( + issues, + 'database.statementTimeoutMs', + config.database?.statementTimeoutMs, + { + min: 1, + }, + ) requireInteger(issues, 'storage.maxUploadSizeMb', config.storage?.maxUploadSizeMb, { min: 1, max: 100, @@ -253,6 +270,18 @@ function requireInteger( } } +function requireOptionalInteger( + issues: RuntimeConfigValidationIssue[], + configPath: string, + value: unknown, + options: { min?: number; max?: number } = {}, +): void { + if (value === undefined || value === null || value === '') { + return + } + requireInteger(issues, configPath, value, options) +} + function isHttpUrl(value: string): boolean { try { const url = new URL(value) diff --git a/apps/backend/src/db/client.ts b/apps/backend/src/db/client.ts index af3df8e2..89dd34ff 100644 --- a/apps/backend/src/db/client.ts +++ b/apps/backend/src/db/client.ts @@ -11,6 +11,9 @@ export function getDb(): Pool { connectionString: String(runtimeConfig.database?.url || '').trim(), ssl: runtimeConfig.database?.ssl ? { rejectUnauthorized: false } : false, max: Number(runtimeConfig.database?.maxConnections || 10), + idleTimeoutMillis: Number(runtimeConfig.database?.idleTimeoutMs || 30_000), + connectionTimeoutMillis: Number(runtimeConfig.database?.connectionTimeoutMs || 5_000), + statement_timeout: Number(runtimeConfig.database?.statementTimeoutMs || 15_000), }) } diff --git a/apps/backend/src/db/migrations/048_database_performance_guards.sql b/apps/backend/src/db/migrations/048_database_performance_guards.sql new file mode 100644 index 00000000..888e9347 --- /dev/null +++ b/apps/backend/src/db/migrations/048_database_performance_guards.sql @@ -0,0 +1,13 @@ +-- 高频打手订单、拼单统计、快手回调扫描和仪表盘查询的索引保护。 +CREATE INDEX IF NOT EXISTS idx_work_order_shares_order_status_include_quantity + ON work_order_shares(work_order_id, status) INCLUDE (quantity); + +CREATE INDEX IF NOT EXISTS idx_work_order_events_order_type + ON work_order_events(work_order_id, event_type); + +CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_status_updated_at + ON fulfillment_tasks(task_status, updated_at DESC); + +CREATE INDEX IF NOT EXISTS idx_kuaishou_industry_vouchers_callback_retry_order + ON kuaishou_industry_vouchers(COALESCE(send_callback_sent_at, created_at), id) + WHERE send_callback_status <> 'success'; diff --git a/apps/backend/src/repositories/realtime-event-repo.ts b/apps/backend/src/repositories/realtime-event-repo.ts index c2c75ac3..31b40d11 100644 --- a/apps/backend/src/repositories/realtime-event-repo.ts +++ b/apps/backend/src/repositories/realtime-event-repo.ts @@ -47,13 +47,15 @@ export async function createRealtimeEvent(input: { const row = result.rows[0] if (!row) throw new Error('realtime_event_create_failed') - // 事件仅用于短期断线补发,控制表体积和重连回放成本。 - await query( - ` - DELETE FROM realtime_events - WHERE id < GREATEST(0, (SELECT COALESCE(MAX(id), 0) - 500 FROM realtime_events)) - `, - ) + // 事件仅用于短期断线补发;每 100 条清理一次,避免每次业务事件都额外触发 DELETE。 + if (Number(row.id) % 100 === 0) { + await query( + ` + DELETE FROM realtime_events + WHERE id < GREATEST(0, (SELECT COALESCE(MAX(id), 0) - 500 FROM realtime_events)) + `, + ) + } return row } diff --git a/apps/backend/src/repositories/worker-platform/work-order-repo.ts b/apps/backend/src/repositories/worker-platform/work-order-repo.ts index 60c8a080..807a1f3d 100644 --- a/apps/backend/src/repositories/worker-platform/work-order-repo.ts +++ b/apps/backend/src/repositories/worker-platform/work-order-repo.ts @@ -26,22 +26,21 @@ import { maybeUpgradeWorkerLevelWithClient } from './worker-repo.js' const WORK_ORDER_SELECT = ` SELECT wo.*, - COALESCE(( - SELECT SUM(wos.quantity) - FROM work_order_shares wos - WHERE wos.work_order_id = wo.id AND wos.status != 'cancelled' - ), 0)::int AS active_share_quantity, - COALESCE(( - SELECT COUNT(*) - FROM work_order_shares wos - WHERE wos.work_order_id = wo.id AND wos.status = 'joined' - ), 0)::int AS pending_share_count, + COALESCE(share_summary.active_share_quantity, 0)::int AS active_share_quantity, + COALESCE(share_summary.pending_share_count, 0)::int AS pending_share_count, wc.name AS category_name, wu.username AS worker_username, wu.display_name AS worker_display_name FROM work_orders wo LEFT JOIN work_categories wc ON wc.id = wo.category_id LEFT JOIN worker_users wu ON wu.id = wo.assigned_worker_id + LEFT JOIN LATERAL ( + SELECT + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, + COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count + FROM work_order_shares wos + WHERE wos.work_order_id = wo.id + ) share_summary ON TRUE ` const WORK_PRODUCT_RULE_SELECT = ` @@ -106,7 +105,22 @@ type WorkOrderStatisticsRow = { } /** 后台按母工单拼单参与情况计算业务状态,open 只代表大厅生命周期。 */ -function adminOperationalStatusExpression(alias = 'wo') { +function adminOperationalStatusExpression(alias = 'wo', shareSummaryAlias = '') { + if (shareSummaryAlias) { + return `CASE + WHEN ${alias}.status = 'open' + AND ${alias}.sharing_enabled IS TRUE + AND COALESCE(${shareSummaryAlias}.active_share_quantity, 0) > 0 + THEN CASE + WHEN COALESCE(${shareSummaryAlias}.active_share_quantity, 0) >= ${alias}.sharing_total_quantity + AND COALESCE(${shareSummaryAlias}.pending_share_count, 0) = 0 + THEN 'pending_acceptance' + ELSE 'in_progress' + END + ELSE ${alias}.status + END` + } + return `CASE WHEN ${alias}.status = 'open' AND ${alias}.sharing_enabled IS TRUE @@ -1324,7 +1338,7 @@ export async function getWorkOrderStatistics(input: ListInput = {}): Promise>'evidenceDueAt', '')::timestamptz > NOW() ` const operationalStatus = input.useOperationalStatus - ? adminOperationalStatusExpression() + ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status' const statusExpression = input.vipEvidencePending ? `CASE WHEN ${vipEvidencePendingWhere} THEN 'vip_evidence_pending' ELSE ${operationalStatus} END` @@ -4357,6 +4371,29 @@ function buildWorkOrderWhere({ let joins = '' let workerSharePlaceholder = '' + if (useOperationalStatus) { + joins += ` + LEFT JOIN ( + SELECT + wos.work_order_id, + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity, + COUNT(*) FILTER (WHERE wos.status = 'joined') AS pending_share_count + FROM work_order_shares wos + GROUP BY wos.work_order_id + ) operational_share_summary ON operational_share_summary.work_order_id = wo.id` + } + + if (excludeFilledSharing) { + joins += ` + LEFT JOIN ( + SELECT + wos.work_order_id, + SUM(wos.quantity) FILTER (WHERE wos.status != 'cancelled') AS active_share_quantity + FROM work_order_shares wos + GROUP BY wos.work_order_id + ) filter_share_summary ON filter_share_summary.work_order_id = wo.id` + } + // 先从打手的整单、历史单和拼单记录组成可见工单集合,再关联工单详情。 // 这样 COUNT 和列表都不会先扫描全量工单,再逐行检查该打手的拼单状态。 if (workerSharingId) { @@ -4403,7 +4440,7 @@ function buildWorkOrderWhere({ if (normalizedStatuses.length > 0) { params.push(normalizedStatuses) filters.push( - `(${useOperationalStatus ? adminOperationalStatusExpression() : 'wo.status'} = ANY($${params.length}::text[]) OR (${vipEvidencePendingWhere}))`, + `(${useOperationalStatus ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status'} = ANY($${params.length}::text[]) OR (${vipEvidencePendingWhere}))`, ) } else { filters.push(`(${vipEvidencePendingWhere})`) @@ -4412,12 +4449,10 @@ function buildWorkOrderWhere({ params.push(normalizedStatuses) const statusesPlaceholder = `$${params.length}` if (workerSharingId) { - filters.push( - `${workerOrderStatusExpression()} = ANY(${statusesPlaceholder}::text[])`, - ) + filters.push(`${workerOrderStatusExpression()} = ANY(${statusesPlaceholder}::text[])`) } else { filters.push( - `${useOperationalStatus ? adminOperationalStatusExpression() : 'wo.status'} = ANY(${statusesPlaceholder}::text[])`, + `${useOperationalStatus ? adminOperationalStatusExpression('wo', 'operational_share_summary') : 'wo.status'} = ANY(${statusesPlaceholder}::text[])`, ) } } @@ -4472,12 +4507,7 @@ function buildWorkOrderWhere({ if (excludeFilledSharing) { filters.push(`( wo.sharing_enabled IS NOT TRUE - OR COALESCE(( - SELECT SUM(wos.quantity) - FROM work_order_shares wos - WHERE wos.work_order_id = wo.id - AND wos.status != 'cancelled' - ), 0) < wo.sharing_total_quantity + OR COALESCE(filter_share_summary.active_share_quantity, 0) < wo.sharing_total_quantity )`) } return { diff --git a/apps/backend/src/services/admin/admin-dashboard-service.ts b/apps/backend/src/services/admin/admin-dashboard-service.ts index 67882967..6aabea8d 100644 --- a/apps/backend/src/services/admin/admin-dashboard-service.ts +++ b/apps/backend/src/services/admin/admin-dashboard-service.ts @@ -4,7 +4,7 @@ import { nowIso } from '../../utils/time.js' import type { JsonObject } from '../../types/json.js' export async function getAdminDashboardSummary() { - const todayPrefix = nowIso().slice(0, 10) + const { todayStartIso, tomorrowStartIso } = resolveChinaTodayRange(nowIso()) const paidPendingClaimStatuses = [ TASK_STATUS.PAID, TASK_STATUS.LINK_GENERATED, @@ -16,31 +16,36 @@ export async function getAdminDashboardSummary() { const result = await query( ` SELECT - (SELECT COUNT(*)::int FROM orders WHERE to_char(created_at AT TIME ZONE 'Asia/Shanghai', 'YYYY-MM-DD') = $1) AS today_orders, ( SELECT COUNT(*)::int - FROM fulfillment_tasks - WHERE task_status = ANY($2::text[]) - ) AS paid_pending_claim, + FROM orders + WHERE created_at >= $1 AND created_at < $2 + ) AS today_orders, ( SELECT COUNT(*)::int FROM fulfillment_tasks WHERE task_status = ANY($3::text[]) + ) AS paid_pending_claim, + ( + SELECT COUNT(*)::int + FROM fulfillment_tasks + WHERE task_status = ANY($4::text[]) ) AS claiming_tasks, ( SELECT COUNT(*)::int FROM fulfillment_tasks - WHERE task_status = $4 - AND to_char(updated_at AT TIME ZONE 'Asia/Shanghai', 'YYYY-MM-DD') = $1 + WHERE task_status = $5 + AND updated_at >= $1 AND updated_at < $2 ) AS redeemed_today, ( SELECT COUNT(*)::int FROM fulfillment_tasks - WHERE task_status = ANY($5::text[]) + WHERE task_status = ANY($6::text[]) ) AS abnormal_tasks `, [ - todayPrefix, + todayStartIso, + tomorrowStartIso, paidPendingClaimStatuses, claimingStatuses, TASK_STATUS.REDEEMED, @@ -57,3 +62,28 @@ export async function getAdminDashboardSummary() { abnormalTasks: Number(summary?.abnormal_tasks || 0), } } + +/** 返回中国时区当天对应的 UTC 查询范围,便于使用时间索引。 */ +function resolveChinaTodayRange(now: string) { + const parts = new Intl.DateTimeFormat('en-US', { + timeZone: 'Asia/Shanghai', + year: 'numeric', + month: '2-digit', + day: '2-digit', + }) + .formatToParts(new Date(now)) + .reduce>((result, part) => { + result[part.type] = part.value + return result + }, {}) + const year = Number(parts.year || 0) + const month = Number(parts.month || 1) + const day = Number(parts.day || 1) + const todayStart = Date.UTC(year || 0, (month || 1) - 1, day || 1, -8) + const tomorrowStart = todayStart + 24 * 60 * 60 * 1000 + + return { + todayStartIso: new Date(todayStart).toISOString(), + tomorrowStartIso: new Date(tomorrowStart).toISOString(), + } +} diff --git a/apps/backend/src/services/worker-platform/worker-service.ts b/apps/backend/src/services/worker-platform/worker-service.ts index 44b0abdb..c51cfabe 100644 --- a/apps/backend/src/services/worker-platform/worker-service.ts +++ b/apps/backend/src/services/worker-platform/worker-service.ts @@ -1389,7 +1389,6 @@ export async function listWorkerHallOrders(query: JsonObject = {}, session: Work const pageSize = normalizePageSize(query.pageSize) const categoryId = normalizeOptionalId(query.categoryId) const worker = await getRequiredWorker(session.workerId) - await settleOverdueWorkOrders({ workerId: worker.id, limit: 50 }) const permissions = resolveWorkerPermissions(worker) const visibleDelaySeconds = Number(permissions.visibleDelaySeconds || 0) const visibleAfterIso = @@ -1650,7 +1649,6 @@ export async function settleOverdueWorkOrders(options: { limit?: number; workerI export async function listWorkerMyOrders(query: JsonObject = {}, session: WorkerSession) { requireActiveWorkerSession(session) const worker = await getRequiredWorker(session.workerId) - await settleOverdueWorkOrders({ workerId: worker.id, limit: 50 }) const page = normalizePage(query.page) const pageSize = normalizePageSize(query.pageSize) const vipEvidencePending = normalizeBoolean(query.vipEvidencePending, false) diff --git a/apps/backend/src/types/runtime-config.ts b/apps/backend/src/types/runtime-config.ts index 3731969d..8cebe280 100644 --- a/apps/backend/src/types/runtime-config.ts +++ b/apps/backend/src/types/runtime-config.ts @@ -34,6 +34,9 @@ export type RuntimeConfig = { url: string ssl: boolean maxConnections: number + idleTimeoutMs?: number + connectionTimeoutMs?: number + statementTimeoutMs?: number } storage: { endpoint: string diff --git a/docker-compose.yml b/docker-compose.yml index 5d78203a..7d497fd1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -52,7 +52,10 @@ services: PORT: ${BACKEND_PORT:-3000} DATABASE_URL: ${DATABASE_URL:-postgres://postgres:postgres@postgres:5432/order_site} DATABASE_SSL: ${DATABASE_SSL:-false} - DATABASE_MAX_CONNECTIONS: ${DATABASE_MAX_CONNECTIONS:-20} + DATABASE_MAX_CONNECTIONS: ${DATABASE_MAX_CONNECTIONS:-8} + DATABASE_IDLE_TIMEOUT_MS: ${DATABASE_IDLE_TIMEOUT_MS:-30000} + DATABASE_CONNECTION_TIMEOUT_MS: ${DATABASE_CONNECTION_TIMEOUT_MS:-5000} + DATABASE_STATEMENT_TIMEOUT_MS: ${DATABASE_STATEMENT_TIMEOUT_MS:-15000} DATA_ROOT: /app/data LOG_RETENTION_DAYS: ${LOG_RETENTION_DAYS:-30} CLAIM_BASE_URL: ${CLAIM_BASE_URL:-http://localhost/#/claim}