fix: concurrent uploads and overwrite-in-place sync strategy
- Upload S3 files 10x concurrently instead of sequentially - Replace clear-then-write with overwrite-in-place + delete stale to avoid S3 404s during concurrent KB ingestion jobs
This commit is contained in:
parent
7240147736
commit
641494530b
1 changed files with 54 additions and 31 deletions
|
|
@ -52,9 +52,10 @@ interface WorkOrderComment {
|
|||
ingested_at?: string;
|
||||
}
|
||||
|
||||
async function clearOldFiles(bucket: string): Promise<void> {
|
||||
/** List all existing S3 keys under the work-orders/ prefix. */
|
||||
async function listExistingKeys(bucket: string): Promise<Set<string>> {
|
||||
const keys = new Set<string>();
|
||||
let continuationToken: string | undefined;
|
||||
const toDelete: { Key: string }[] = [];
|
||||
|
||||
do {
|
||||
const res = await s3.send(
|
||||
|
|
@ -65,23 +66,27 @@ async function clearOldFiles(bucket: string): Promise<void> {
|
|||
}),
|
||||
);
|
||||
for (const obj of res.Contents ?? []) {
|
||||
if (obj.Key) toDelete.push({ Key: obj.Key });
|
||||
if (obj.Key) keys.add(obj.Key);
|
||||
}
|
||||
continuationToken = res.NextContinuationToken;
|
||||
} while (continuationToken);
|
||||
|
||||
if (toDelete.length === 0) return;
|
||||
return keys;
|
||||
}
|
||||
|
||||
// S3 DeleteObjects accepts up to 1000 keys per request
|
||||
for (let i = 0; i < toDelete.length; i += 1000) {
|
||||
/** Delete S3 keys that no longer correspond to active records. */
|
||||
async function deleteStaleKeys(bucket: string, keys: string[]): Promise<void> {
|
||||
if (keys.length === 0) return;
|
||||
|
||||
for (let i = 0; i < keys.length; i += 1000) {
|
||||
await s3.send(
|
||||
new DeleteObjectsCommand({
|
||||
Bucket: bucket,
|
||||
Delete: { Objects: toDelete.slice(i, i + 1000) },
|
||||
Delete: { Objects: keys.slice(i, i + 1000).map((Key) => ({ Key })) },
|
||||
}),
|
||||
);
|
||||
}
|
||||
console.log(`Deleted ${toDelete.length} stale work order file(s) from S3`);
|
||||
console.log(`Deleted ${keys.length} stale work order file(s) from S3`);
|
||||
}
|
||||
|
||||
async function scanAllWorkOrders(): Promise<WorkOrder[]> {
|
||||
|
|
@ -194,44 +199,62 @@ export const handler = async (): Promise<void> => {
|
|||
const kbId = process.env.KNOWLEDGE_BASE_ID!;
|
||||
const dsId = process.env.DATA_SOURCE_ID!;
|
||||
|
||||
console.log('Clearing stale work order files from S3...');
|
||||
await clearOldFiles(bucket);
|
||||
console.log('Listing existing work order files in S3...');
|
||||
const existingKeys = await listExistingKeys(bucket);
|
||||
console.log(`Found ${existingKeys.size} existing file(s) in S3`);
|
||||
|
||||
console.log('Scanning work orders from DynamoDB...');
|
||||
const workOrders = await scanAllWorkOrders();
|
||||
console.log(`Found ${workOrders.length} work order(s)`);
|
||||
|
||||
const writtenKeys = new Set<string>();
|
||||
let synced = 0;
|
||||
let failed = 0;
|
||||
|
||||
for (const wo of workOrders) {
|
||||
try {
|
||||
const comments = await getComments(wo.work_order_id);
|
||||
const markdown = formatWorkOrderMarkdown(wo, comments);
|
||||
// Upload in batches of 10 (each WO also queries comments, so keep concurrency moderate)
|
||||
const CONCURRENCY = 10;
|
||||
for (let i = 0; i < workOrders.length; i += CONCURRENCY) {
|
||||
const batch = workOrders.slice(i, i + CONCURRENCY);
|
||||
const results = await Promise.allSettled(
|
||||
batch.map(async (wo) => {
|
||||
const comments = await getComments(wo.work_order_id);
|
||||
const markdown = formatWorkOrderMarkdown(wo, comments);
|
||||
const key = `${WO_PREFIX}${wo.work_order_id}.md`;
|
||||
|
||||
await s3.send(
|
||||
new PutObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: `${WO_PREFIX}${wo.work_order_id}.md`,
|
||||
Body: markdown,
|
||||
ContentType: 'text/markdown',
|
||||
Metadata: {
|
||||
'work-order-id': wo.work_order_id,
|
||||
'site-code': (wo.site_code ?? '').slice(0, 256),
|
||||
},
|
||||
}),
|
||||
);
|
||||
await s3.send(
|
||||
new PutObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: key,
|
||||
Body: markdown,
|
||||
ContentType: 'text/markdown',
|
||||
Metadata: {
|
||||
'work-order-id': wo.work_order_id,
|
||||
'site-code': (wo.site_code ?? '').slice(0, 256),
|
||||
},
|
||||
}),
|
||||
);
|
||||
|
||||
console.log(`Synced: ${wo.work_order_id} (${comments.length} comment(s))`);
|
||||
synced++;
|
||||
} catch (err) {
|
||||
console.error(`Failed to sync work order ${wo.work_order_id}:`, err);
|
||||
failed++;
|
||||
return key;
|
||||
}),
|
||||
);
|
||||
|
||||
for (const r of results) {
|
||||
if (r.status === 'fulfilled') {
|
||||
writtenKeys.add(r.value);
|
||||
synced++;
|
||||
} else {
|
||||
console.error(`Failed to sync work order:`, r.reason);
|
||||
failed++;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
console.log(`Upload complete. ${synced} synced, ${failed} failed.`);
|
||||
|
||||
// Delete S3 files for records that no longer exist in DynamoDB
|
||||
const staleKeys = [...existingKeys].filter((k) => !writtenKeys.has(k));
|
||||
await deleteStaleKeys(bucket, staleKeys);
|
||||
|
||||
console.log('Triggering Bedrock KB ingestion job...');
|
||||
const ingestionRes = await bedrockAgent.send(
|
||||
new StartIngestionJobCommand({
|
||||
|
|
|
|||
Reference in a new issue