perf(import): parallelize ID allocation and raise clone concurrency to 10
Replace the serialized SELECT MAX + INSERT id-allocation chains for items_base and catalog_items with a lazy-seeded in-process counter so concurrent clone workers no longer queue on a global lock per item. Re-seeds after 60s idle to avoid colliding with externally added rows. Raise the clone import concurrency default from 6 to the batch cap of 10.
This commit is contained in:
1 parent
742536820a
commit
172941c542
8 files changed
+144
-52
No files matched your search
@@ -454,7 +454,7 @@ function FurniGrid({ source }: FurniGridProps) {
|
||||
{
|
||||
sourceId: source.id,
|
||||
classnames: toClone.map((it) => it.classname),
|
||||
concurrency: 6,
|
||||
concurrency: 10,
|
||||
},
|
||||
(classname) => {
|
||||
done++;
|
||||
@@ -545,7 +545,7 @@ function FurniGrid({ source }: FurniGridProps) {
|
||||
const chunk = names.slice(i, i + CHUNK);
|
||||
await runSseImport(
|
||||
"/api/admin/import/clone/batch",
|
||||
{ sourceId: source.id, classnames: chunk, concurrency: 6 },
|
||||
{ sourceId: source.id, classnames: chunk, concurrency: 10 },
|
||||
(classname) => {
|
||||
done++;
|
||||
markDone(classname);
|
||||
|
||||
@@ -67,7 +67,7 @@ export const POST = withAdmin(
|
||||
|
||||
return runSseBatch<BatchItem>({
|
||||
items,
|
||||
concurrency: body.concurrency || 6,
|
||||
concurrency: body.concurrency || 10,
|
||||
signal: request.signal,
|
||||
labelOf: (it) => it.classname,
|
||||
worker: async (it, _index, report) => {
|
||||
|
||||
@@ -108,7 +108,7 @@ export const POST = withAdmin(
|
||||
|
||||
return runSseBatch({
|
||||
items: allItems,
|
||||
concurrency: 6,
|
||||
concurrency: 10,
|
||||
signal: request.signal,
|
||||
labelOf: (it) => `${it.sourceName}/${it.classname}`,
|
||||
worker: async (it, _index, report) => {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// @vitest-environment node
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
// biome-ignore lint/suspicious/noExplicitAny: test helper for mocking with arbitrary signatures
|
||||
type AnyFn = (...args: any[]) => any;
|
||||
@@ -69,7 +69,11 @@ vi.mock("@/lib/db", () => ({
|
||||
ItemsBase: { id: "id", itemName: "itemName" },
|
||||
}));
|
||||
|
||||
import { cloneSingleFurni, parseFurnidata } from "./clone-import";
|
||||
import {
|
||||
__resetItemsBaseAllocatorForTests,
|
||||
cloneSingleFurni,
|
||||
parseFurnidata,
|
||||
} from "./clone-import";
|
||||
|
||||
const SOURCE = {
|
||||
id: "s",
|
||||
@@ -92,6 +96,10 @@ const ENTRY = {
|
||||
};
|
||||
|
||||
describe("clone-import", () => {
|
||||
beforeEach(() => {
|
||||
__resetItemsBaseAllocatorForTests();
|
||||
});
|
||||
|
||||
it("parseFurnidata normalizes room + wall items with itemType", () => {
|
||||
const list = parseFurnidata({
|
||||
roomitemtypes: { furnitype: [ENTRY] },
|
||||
|
||||
@@ -137,26 +137,54 @@ export async function fetchSourceFurnidata(
|
||||
* concurrent clones in the same Node process cannot race on MAX(id)+1.
|
||||
* Mirrors the `allocateCatalogItemId` pattern from furni-import.ts.
|
||||
*/
|
||||
let itemsBaseIdAllocChain: Promise<unknown> = Promise.resolve();
|
||||
let itemsBaseNextId: number | null = null;
|
||||
let itemsBaseIdLastUsed = 0;
|
||||
let itemsBaseIdSeedChain: Promise<void> = Promise.resolve();
|
||||
const ITEMS_BASE_ID_REFRESH_MS = 60_000;
|
||||
|
||||
/** Allocate a unique items_base id. Seeded once with MAX(id)+1, then handed
|
||||
* out from an in-process counter so concurrent clones don't serialize on a
|
||||
* SELECT MAX + INSERT round-trip per item. Re-seeds when idle for a while so
|
||||
* rows added by the emulator or other processes don't collide. */
|
||||
async function allocateItemsBaseId<T>(
|
||||
insertFn: (nextId: number) => Promise<T>,
|
||||
): Promise<T> {
|
||||
const prev = itemsBaseIdAllocChain;
|
||||
let settle!: () => void;
|
||||
itemsBaseIdAllocChain = new Promise<void>((r) => {
|
||||
settle = r;
|
||||
});
|
||||
await prev.catch(() => {});
|
||||
try {
|
||||
const [idRows] = (await db.execute(sql`
|
||||
SELECT COALESCE(MAX(id), 0) + 1 AS next FROM items_base
|
||||
`)) as unknown as [Array<{ next: number }>, unknown];
|
||||
const nextId = Number(idRows[0]?.next ?? 1);
|
||||
return await insertFn(nextId);
|
||||
} finally {
|
||||
settle();
|
||||
const now = Date.now();
|
||||
if (
|
||||
itemsBaseNextId === null ||
|
||||
now - itemsBaseIdLastUsed > ITEMS_BASE_ID_REFRESH_MS
|
||||
) {
|
||||
let settle!: () => void;
|
||||
const prev = itemsBaseIdSeedChain;
|
||||
itemsBaseIdSeedChain = new Promise<void>((r) => {
|
||||
settle = r;
|
||||
});
|
||||
await prev.catch(() => {});
|
||||
try {
|
||||
if (
|
||||
itemsBaseNextId === null ||
|
||||
Date.now() - itemsBaseIdLastUsed > ITEMS_BASE_ID_REFRESH_MS
|
||||
) {
|
||||
const [idRows] = (await db.execute(sql`
|
||||
SELECT COALESCE(MAX(id), 0) + 1 AS next FROM items_base
|
||||
`)) as unknown as [Array<{ next: number }>, unknown];
|
||||
itemsBaseNextId = Number(idRows[0]?.next ?? 1);
|
||||
}
|
||||
} finally {
|
||||
settle();
|
||||
}
|
||||
}
|
||||
itemsBaseIdLastUsed = Date.now();
|
||||
const nextId = itemsBaseNextId;
|
||||
itemsBaseNextId += 1;
|
||||
return await insertFn(nextId);
|
||||
}
|
||||
|
||||
/** Test seam: drop the cached id counter between tests. */
|
||||
export function __resetItemsBaseAllocatorForTests(): void {
|
||||
itemsBaseNextId = null;
|
||||
itemsBaseIdLastUsed = 0;
|
||||
itemsBaseIdSeedChain = Promise.resolve();
|
||||
}
|
||||
|
||||
export async function cloneSingleFurni(params: {
|
||||
|
||||
@@ -276,26 +276,54 @@ export function autoPriceFurni(classname: string): {
|
||||
* connection for the whole operation; the in-process mutex is enough
|
||||
* for the single-server deployment and avoids pool-connection gymnastics.
|
||||
*/
|
||||
let catalogIdAllocChain: Promise<unknown> = Promise.resolve();
|
||||
let catalogNextId: number | null = null;
|
||||
let catalogIdLastUsed = 0;
|
||||
let catalogIdSeedChain: Promise<void> = Promise.resolve();
|
||||
const CATALOG_ID_REFRESH_MS = 60_000;
|
||||
|
||||
/** Allocate a unique catalog_items id. Seeded once with MAX(id)+1, then handed
|
||||
* out from an in-process counter so concurrent imports don't serialize on a
|
||||
* SELECT MAX + INSERT round-trip per item. Re-seeds when idle so rows added
|
||||
* externally don't collide. */
|
||||
export async function allocateCatalogItemId<T>(
|
||||
insertFn: (nextId: number) => Promise<T>,
|
||||
): Promise<T> {
|
||||
const prev = catalogIdAllocChain;
|
||||
let settle!: () => void;
|
||||
catalogIdAllocChain = new Promise<void>((r) => {
|
||||
settle = r;
|
||||
});
|
||||
await prev.catch(() => {});
|
||||
try {
|
||||
const [maxIdResult] = (await db.execute(sql`
|
||||
SELECT MAX(id) as maxId FROM catalog_items
|
||||
`)) as unknown as [Array<{ maxId: number | bigint | null }>, unknown];
|
||||
const nextId = Number(maxIdResult[0]?.maxId ?? 0) + 1;
|
||||
return await insertFn(nextId);
|
||||
} finally {
|
||||
settle();
|
||||
const now = Date.now();
|
||||
if (
|
||||
catalogNextId === null ||
|
||||
now - catalogIdLastUsed > CATALOG_ID_REFRESH_MS
|
||||
) {
|
||||
let settle!: () => void;
|
||||
const prev = catalogIdSeedChain;
|
||||
catalogIdSeedChain = new Promise<void>((r) => {
|
||||
settle = r;
|
||||
});
|
||||
await prev.catch(() => {});
|
||||
try {
|
||||
if (
|
||||
catalogNextId === null ||
|
||||
Date.now() - catalogIdLastUsed > CATALOG_ID_REFRESH_MS
|
||||
) {
|
||||
const [maxIdResult] = (await db.execute(sql`
|
||||
SELECT MAX(id) as maxId FROM catalog_items
|
||||
`)) as unknown as [Array<{ maxId: number | bigint | null }>, unknown];
|
||||
catalogNextId = Number(maxIdResult[0]?.maxId ?? 0) + 1;
|
||||
}
|
||||
} finally {
|
||||
settle();
|
||||
}
|
||||
}
|
||||
catalogIdLastUsed = Date.now();
|
||||
const nextId = catalogNextId;
|
||||
catalogNextId += 1;
|
||||
return await insertFn(nextId);
|
||||
}
|
||||
|
||||
/** Test seam: drop the cached id counter between tests. */
|
||||
export function __resetCatalogIdAllocatorForTests(): void {
|
||||
catalogNextId = null;
|
||||
catalogIdLastUsed = 0;
|
||||
catalogIdSeedChain = Promise.resolve();
|
||||
}
|
||||
|
||||
/** Get or create a category sub-page under the imported parent page. */
|
||||
|
||||
@@ -82,12 +82,16 @@ vi.mock("@/lib/services/swf/nitro-builder", () => ({
|
||||
parseNitroBundle: vi.fn(() => ({ json: { xdim: 1, ydim: 1 } })),
|
||||
}));
|
||||
|
||||
import { uploadSingleFurni } from "./upload-import";
|
||||
import {
|
||||
__resetUploadItemsBaseAllocatorForTests,
|
||||
uploadSingleFurni,
|
||||
} from "./upload-import";
|
||||
|
||||
describe("uploadSingleFurni live asset mirrors", () => {
|
||||
let tempDir: string;
|
||||
|
||||
beforeEach(async () => {
|
||||
__resetUploadItemsBaseAllocatorForTests();
|
||||
tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "upload-furni-"));
|
||||
assetTargets.swfDir = path.join(tempDir, "cms", "swf");
|
||||
assetTargets.iconDir = path.join(tempDir, "cms", "icons");
|
||||
|
||||
@@ -41,7 +41,10 @@ const MIGRATIONS_DIR = path.resolve(
|
||||
"drizzle/migrations",
|
||||
);
|
||||
|
||||
let itemsBaseIdAllocChain: Promise<unknown> = Promise.resolve();
|
||||
let itemsBaseNextId: number | null = null;
|
||||
let itemsBaseIdLastUsed = 0;
|
||||
let itemsBaseIdSeedChain: Promise<void> = Promise.resolve();
|
||||
const ITEMS_BASE_ID_REFRESH_MS = 60_000;
|
||||
|
||||
function pathKey(value: string): string {
|
||||
const normalized = path.normalize(value);
|
||||
@@ -76,21 +79,42 @@ async function mirrorUploadedAsset(
|
||||
async function allocateItemsBaseId<T>(
|
||||
insertFn: (nextId: number) => Promise<T>,
|
||||
): Promise<T> {
|
||||
const prev = itemsBaseIdAllocChain;
|
||||
let settle!: () => void;
|
||||
itemsBaseIdAllocChain = new Promise<void>((r) => {
|
||||
settle = r;
|
||||
});
|
||||
await prev.catch(() => {});
|
||||
try {
|
||||
const [idRows] = (await db.execute(sql`
|
||||
SELECT COALESCE(MAX(id), 0) + 1 AS next FROM items_base
|
||||
`)) as unknown as [Array<{ next: number }>, unknown];
|
||||
const nextId = Number(idRows[0]?.next ?? 1);
|
||||
return await insertFn(nextId);
|
||||
} finally {
|
||||
settle();
|
||||
const now = Date.now();
|
||||
if (
|
||||
itemsBaseNextId === null ||
|
||||
now - itemsBaseIdLastUsed > ITEMS_BASE_ID_REFRESH_MS
|
||||
) {
|
||||
let settle!: () => void;
|
||||
const prev = itemsBaseIdSeedChain;
|
||||
itemsBaseIdSeedChain = new Promise<void>((r) => {
|
||||
settle = r;
|
||||
});
|
||||
await prev.catch(() => {});
|
||||
try {
|
||||
if (
|
||||
itemsBaseNextId === null ||
|
||||
Date.now() - itemsBaseIdLastUsed > ITEMS_BASE_ID_REFRESH_MS
|
||||
) {
|
||||
const [idRows] = (await db.execute(sql`
|
||||
SELECT COALESCE(MAX(id), 0) + 1 AS next FROM items_base
|
||||
`)) as unknown as [Array<{ next: number }>, unknown];
|
||||
itemsBaseNextId = Number(idRows[0]?.next ?? 1);
|
||||
}
|
||||
} finally {
|
||||
settle();
|
||||
}
|
||||
}
|
||||
itemsBaseIdLastUsed = Date.now();
|
||||
const nextId = itemsBaseNextId;
|
||||
itemsBaseNextId += 1;
|
||||
return await insertFn(nextId);
|
||||
}
|
||||
|
||||
/** Test seam: drop the cached id counter between tests. */
|
||||
export function __resetUploadItemsBaseAllocatorForTests(): void {
|
||||
itemsBaseNextId = null;
|
||||
itemsBaseIdLastUsed = 0;
|
||||
itemsBaseIdSeedChain = Promise.resolve();
|
||||
}
|
||||
|
||||
function escapeSql(val: string | number): string {
|
||||
|
||||
Reference in new issue
Block a user