perf: reduce worker order list query load

This commit is contained in:
yml2213
2026-08-30 16:40:09 +08:00
parent a45d7597fb
commit d24c69141a
8 changed files with 223 additions and 117 deletions
@@ -110,6 +110,12 @@ test('后台接单工单支持按接单分类筛选', () => {
assert.deepEqual(query.params, [7])
})
test('空白搜索词不生成全表 LIKE 条件', () => {
const query = buildWorkOrderWhere({ keyword: ' ' })
assert.doesNotMatch(query.whereClause, /LIKE LOWER/)
assert.deepEqual(query.params, [])
})
test('大厅候选按分类上限取队首后再受全局上限限制', () => {
const query = buildWorkOrderWhere({
hallCandidateLimit: 20,
@@ -58,6 +58,13 @@ export const WORK_ORDER_SELECT = `
) share_summary ON TRUE
`
// Keep the page and total in one filtered query. This avoids running the same
// joins and share aggregation twice for every list request.
const LIST_WORK_ORDER_SELECT = WORK_ORDER_SELECT.replace(
' SELECT\n',
' SELECT\n COUNT(*) OVER()::int AS total_count,\n',
)
type WorkerOrderOverviewRow = {
status: string
order_count: number
@@ -268,15 +275,6 @@ export async function listWorkOrders({
createdFrom,
createdTo,
})
const totalResult = await query<{ total: number }>(
`
SELECT COUNT(*)::int AS total
FROM work_orders wo
${joins}
${whereClause}
`,
params,
)
const offset = (page - 1) * pageSize
const orderBy =
sort === 'worker_claimed_at_desc' && workerSharingId
@@ -292,14 +290,30 @@ export async function listWorkOrders({
const effectiveOrderBy = pinnedFirst ? `wo.pinned_at DESC NULLS LAST, ${orderBy}` : orderBy
params.push(pageSize, offset)
const itemsResult = await query<WorkOrderRow>(
`${WORK_ORDER_SELECT}
`${LIST_WORK_ORDER_SELECT}
${joins}
${whereClause}
ORDER BY ${effectiveOrderBy}
LIMIT $${params.length - 1} OFFSET $${params.length}`,
params,
)
return { items: itemsResult.rows, total: Number(totalResult.rows[0]?.total || 0) }
// A page beyond the end has no row from which to read the window count.
// Keep that uncommon pagination edge case correct without restoring a
// second query for normal requests.
let total = Number(
(itemsResult.rows[0] as WorkOrderRow & { total_count?: number })?.total_count || 0,
)
if (itemsResult.rows.length === 0 && page > 1) {
const totalResult = await query<{ total: number }>(
`SELECT COUNT(*)::int AS total FROM work_orders wo ${joins} ${whereClause}`,
params.slice(0, -2),
)
total = Number(totalResult.rows[0]?.total || 0)
}
return {
items: itemsResult.rows,
total,
}
}
/** 按筛选条件汇总工单,统计范围不受当前分页影响。 */
@@ -20,6 +20,10 @@ import {
requireActiveWorkerSession,
type WorkerSession,
} from './worker-session-context-service.js'
import { WorkerListCache } from './worker-list-cache.js'
const HALL_LIST_CACHE_TTL_MS = 2_000
const hallListCache = new WorkerListCache(HALL_LIST_CACHE_TTL_MS)
export async function listWorkerHallOrders(query: JsonObject = {}, session: WorkerSession) {
requireActiveWorkerSession(session)
@@ -29,62 +33,66 @@ export async function listWorkerHallOrders(query: JsonObject = {}, session: Work
const worker = await getRequiredWorker(session.workerId)
const permissions = resolveWorkerPermissions(worker)
const visibleDelaySeconds = Number(permissions.visibleDelaySeconds || 0)
const visibleAfterIso =
visibleDelaySeconds > 0 ? new Date(Date.now() - visibleDelaySeconds * 1000).toISOString() : ''
const hallConfig = getWorkerHallConfig()
const { items, total } = await listWorkOrders({
page,
pageSize,
statuses: [WORK_ORDER_STATUS.OPEN],
keyword: String(query.keyword || '').trim(),
categoryId: categoryId || 0,
hallCandidateLimit: hallConfig.maxVisibleOrders,
hallCategoryLimits: hallConfig.categoryMaxVisibleOrders,
visibleAfterIso,
excludeFilledSharing: true,
pinnedFirst: true,
})
const categories = (await listWorkCategories()).map(mapWorkCategory)
const orderIds = items.map((item) => Number(item.id))
const shares = await listWorkOrderSharesByOrderIds(orderIds)
const myShareByOrderId = new Map(
shares
.filter(
(share) => share.status !== 'cancelled' && Number(share.worker_id) === Number(worker.id),
)
.map((share) => [Number(share.work_order_id), share]),
)
const sharingProgressByOrderId = new Map<
number,
{ joinedQuantity: number; pendingSubmissionCount: number }
>()
for (const share of shares) {
if (share.status === 'cancelled') continue
const orderId = Number(share.work_order_id)
const progress = sharingProgressByOrderId.get(orderId) || {
joinedQuantity: 0,
pendingSubmissionCount: 0,
}
progress.joinedQuantity += Number(share.quantity || 0)
if (share.status === 'joined') progress.pendingSubmissionCount += 1
sharingProgressByOrderId.set(orderId, progress)
}
return {
items: items.map((item) => {
const myShare = myShareByOrderId.get(Number(item.id))
return {
...mapWorkOrderForWorker(item, permissions),
myShare: myShare ? mapWorkOrderShare(myShare) : null,
sharingProgress: sharingProgressByOrderId.get(Number(item.id)) || {
joinedQuantity: 0,
pendingSubmissionCount: 0,
},
const keyword = String(query.keyword || '').trim()
const cacheKey = JSON.stringify([worker.id, page, pageSize, keyword, categoryId || 0])
return hallListCache.getOrSet(cacheKey, async () => {
const visibleAfterIso =
visibleDelaySeconds > 0 ? new Date(Date.now() - visibleDelaySeconds * 1000).toISOString() : ''
const { items, total } = await listWorkOrders({
page,
pageSize,
statuses: [WORK_ORDER_STATUS.OPEN],
keyword,
categoryId: categoryId || 0,
hallCandidateLimit: hallConfig.maxVisibleOrders,
hallCategoryLimits: hallConfig.categoryMaxVisibleOrders,
visibleAfterIso,
excludeFilledSharing: true,
pinnedFirst: true,
})
const categories = (await listWorkCategories()).map(mapWorkCategory)
const orderIds = items.map((item) => Number(item.id))
const shares = await listWorkOrderSharesByOrderIds(orderIds)
const myShareByOrderId = new Map(
shares
.filter(
(share) => share.status !== 'cancelled' && Number(share.worker_id) === Number(worker.id),
)
.map((share) => [Number(share.work_order_id), share]),
)
const sharingProgressByOrderId = new Map<
number,
{ joinedQuantity: number; pendingSubmissionCount: number }
>()
for (const share of shares) {
if (share.status === 'cancelled') continue
const orderId = Number(share.work_order_id)
const progress = sharingProgressByOrderId.get(orderId) || {
joinedQuantity: 0,
pendingSubmissionCount: 0,
}
}),
pagination: { page, pageSize, total },
categories,
visibleDelaySeconds,
}
progress.joinedQuantity += Number(share.quantity || 0)
if (share.status === 'joined') progress.pendingSubmissionCount += 1
sharingProgressByOrderId.set(orderId, progress)
}
return {
items: items.map((item) => {
const myShare = myShareByOrderId.get(Number(item.id))
return {
...mapWorkOrderForWorker(item, permissions),
myShare: myShare ? mapWorkOrderShare(myShare) : null,
sharingProgress: sharingProgressByOrderId.get(Number(item.id)) || {
joinedQuantity: 0,
pendingSubmissionCount: 0,
},
}
}),
pagination: { page, pageSize, total },
categories,
visibleDelaySeconds,
}
})
}
export async function getWorkerHallLeaderboard(session: WorkerSession, period: unknown = 'all') {
@@ -0,0 +1,26 @@
import assert from 'node:assert/strict'
import test from 'node:test'
import { WorkerListCache } from './worker-list-cache.js'
test('WorkerListCache coalesces concurrent reads and expires entries', async () => {
const cache = new WorkerListCache(15)
let calls = 0
const factory = async () => {
calls += 1
await new Promise((resolve) => setTimeout(resolve, 1))
return { calls }
}
const [first, second] = await Promise.all([
cache.getOrSet('same', factory),
cache.getOrSet('same', factory),
])
assert.deepEqual(first, { calls: 1 })
assert.deepEqual(second, { calls: 1 })
assert.equal(calls, 1)
await new Promise((resolve) => setTimeout(resolve, 20))
assert.deepEqual(await cache.getOrSet('same', factory), { calls: 2 })
assert.equal(calls, 2)
})
@@ -0,0 +1,32 @@
type CacheEntry = {
expiresAt: number
value: Promise<unknown>
}
/** Short-lived read cache that also coalesces concurrent requests. */
export class WorkerListCache {
private readonly entries = new Map<string, CacheEntry>()
constructor(
private readonly ttlMs: number,
private readonly maxEntries = 256,
) {}
getOrSet<T>(key: string, factory: () => Promise<T>): Promise<T> {
const now = Date.now()
const cached = this.entries.get(key)
if (cached && cached.expiresAt > now) return cached.value as Promise<T>
if (cached) this.entries.delete(key)
const value = factory().catch((error) => {
this.entries.delete(key)
throw error
})
if (this.entries.size >= this.maxEntries) {
const oldestKey = this.entries.keys().next().value
if (oldestKey) this.entries.delete(oldestKey)
}
this.entries.set(key, { expiresAt: now + this.ttlMs, value })
return value
}
}
@@ -44,6 +44,10 @@ import {
type WorkerSession,
} from './worker-session-context-service.js'
import { settleOverdueWorkOrders } from './worker-order-timeout-service.js'
import { WorkerListCache } from './worker-list-cache.js'
const MY_ORDER_LIST_CACHE_TTL_MS = 2_000
const myOrderListCache = new WorkerListCache(MY_ORDER_LIST_CACHE_TTL_MS)
export async function listWorkerMyOrders(query: JsonObject = {}, session: WorkerSession) {
requireActiveWorkerSession(session)
@@ -51,62 +55,74 @@ export async function listWorkerMyOrders(query: JsonObject = {}, session: Worker
const page = normalizePage(query.page)
const pageSize = normalizePageSize(query.pageSize)
const vipEvidencePending = normalizeBoolean(query.vipEvidencePending, false)
const { items, total } = await listWorkOrders({
const permissions = resolveWorkerPermissions(worker)
const statuses = normalizeStatuses(query.status)
const keyword = String(query.keyword || '').trim()
const cacheKey = JSON.stringify([
worker.id,
page,
pageSize,
statuses: normalizeStatuses(query.status),
keyword: String(query.keyword || '').trim(),
workerSharingId: worker.id,
statuses,
keyword,
vipEvidencePending,
sort: 'worker_claimed_at_desc',
])
return myOrderListCache.getOrSet(cacheKey, async () => {
const { items, total } = await listWorkOrders({
page,
pageSize,
statuses,
keyword,
workerSharingId: worker.id,
vipEvidencePending,
sort: 'worker_claimed_at_desc',
})
const workOrderIds = items.map((item) => Number(item.id))
const myShares = await listWorkerSharesByWorker(worker.id, workOrderIds)
const myShareByOrderId = new Map(myShares.map((share) => [Number(share.work_order_id), share]))
const [notes, views, pendingCancelOrderIds, pendingFeedbackOrderIds, materialEventAtByOrderId] =
await Promise.all([
listWorkerWorkOrderNotes(worker.id, workOrderIds),
listWorkerWorkOrderViews(worker.id, workOrderIds),
getPendingCancelRequestOrderIds(worker.id, workOrderIds),
getPendingFeedbackOrderIds(worker.id, workOrderIds),
getLatestMaterialEventAtByOrderIds(workOrderIds),
])
const noteByOrderId = new Map(
notes.map((item) => [Number(item.work_order_id), String(item.note || '')]),
)
const viewedAtByOrderId = new Map(
views.map((item) => [Number(item.work_order_id), String(item.viewed_at || '')]),
)
const materialSeenAtByOrderId = new Map(
views.map((item) => [Number(item.work_order_id), String(item.material_seen_at || '')]),
)
return {
items: items.map((item) => {
const workOrderId = Number(item.id)
const materialEventAt = materialEventAtByOrderId.get(workOrderId) || ''
const materialSeenAt = materialSeenAtByOrderId.get(workOrderId) || ''
const materialUpdated =
Boolean(materialEventAt) &&
(!materialSeenAt ||
new Date(materialEventAt).getTime() > new Date(materialSeenAt).getTime())
return {
...mapWorkOrderForWorker(
item,
permissions,
myShareByOrderId.get(workOrderId),
noteByOrderId.get(workOrderId) || '',
viewedAtByOrderId.get(workOrderId) || null,
),
// 问题单分组内的来源角标:平台标记 / 退单审核中 / 反馈处理中。
problemMarked: item.status === 'problem',
pendingCancelRequest: pendingCancelOrderIds.has(workOrderId),
pendingFeedback: pendingFeedbackOrderIds.has(workOrderId),
materialUpdated,
}
}),
pagination: { page, pageSize, total },
}
})
const permissions = resolveWorkerPermissions(worker)
const workOrderIds = items.map((item) => Number(item.id))
const myShares = await listWorkerSharesByWorker(worker.id, workOrderIds)
const myShareByOrderId = new Map(myShares.map((share) => [Number(share.work_order_id), share]))
const [notes, views, pendingCancelOrderIds, pendingFeedbackOrderIds, materialEventAtByOrderId] =
await Promise.all([
listWorkerWorkOrderNotes(worker.id, workOrderIds),
listWorkerWorkOrderViews(worker.id, workOrderIds),
getPendingCancelRequestOrderIds(worker.id, workOrderIds),
getPendingFeedbackOrderIds(worker.id, workOrderIds),
getLatestMaterialEventAtByOrderIds(workOrderIds),
])
const noteByOrderId = new Map(
notes.map((item) => [Number(item.work_order_id), String(item.note || '')]),
)
const viewedAtByOrderId = new Map(
views.map((item) => [Number(item.work_order_id), String(item.viewed_at || '')]),
)
const materialSeenAtByOrderId = new Map(
views.map((item) => [Number(item.work_order_id), String(item.material_seen_at || '')]),
)
return {
items: items.map((item) => {
const workOrderId = Number(item.id)
const materialEventAt = materialEventAtByOrderId.get(workOrderId) || ''
const materialSeenAt = materialSeenAtByOrderId.get(workOrderId) || ''
const materialUpdated =
Boolean(materialEventAt) &&
(!materialSeenAt ||
new Date(materialEventAt).getTime() > new Date(materialSeenAt).getTime())
return {
...mapWorkOrderForWorker(
item,
permissions,
myShareByOrderId.get(workOrderId),
noteByOrderId.get(workOrderId) || '',
viewedAtByOrderId.get(workOrderId) || null,
),
// 问题单分组内的来源角标:平台标记 / 退单审核中 / 反馈处理中。
problemMarked: item.status === 'problem',
pendingCancelRequest: pendingCancelOrderIds.has(workOrderId),
pendingFeedback: pendingFeedbackOrderIds.has(workOrderId),
materialUpdated,
}
}),
pagination: { page, pageSize, total },
}
}
/** 打手订单标签概览:订单数与预计或已结算金额。 */