fix: concurrent uploads, overwrite-in-place, 15min timeout
- Bump Lambda timeout from 5min to 15min (9k+ POs need more time) - Upload S3 files 25x 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
615970eba2
commit
adbe19aa16
2 changed files with 57 additions and 35 deletions
|
|
@ -133,9 +133,10 @@ async function scanAllPOs(): Promise<Record<string, any>[]> {
|
||||||
return records;
|
return records;
|
||||||
}
|
}
|
||||||
|
|
||||||
async function clearOldFiles(bucket: string): Promise<void> {
|
/** List all existing S3 keys under the purchase-orders/ prefix. */
|
||||||
|
async function listExistingKeys(bucket: string): Promise<Set<string>> {
|
||||||
|
const keys = new Set<string>();
|
||||||
let continuationToken: string | undefined;
|
let continuationToken: string | undefined;
|
||||||
const toDelete: { Key: string }[] = [];
|
|
||||||
|
|
||||||
do {
|
do {
|
||||||
const res = await s3.send(
|
const res = await s3.send(
|
||||||
|
|
@ -146,22 +147,27 @@ async function clearOldFiles(bucket: string): Promise<void> {
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
for (const obj of res.Contents ?? []) {
|
for (const obj of res.Contents ?? []) {
|
||||||
if (obj.Key) toDelete.push({ Key: obj.Key });
|
if (obj.Key) keys.add(obj.Key);
|
||||||
}
|
}
|
||||||
continuationToken = res.NextContinuationToken;
|
continuationToken = res.NextContinuationToken;
|
||||||
} while (continuationToken);
|
} while (continuationToken);
|
||||||
|
|
||||||
if (toDelete.length === 0) return;
|
return keys;
|
||||||
|
}
|
||||||
|
|
||||||
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(
|
await s3.send(
|
||||||
new DeleteObjectsCommand({
|
new DeleteObjectsCommand({
|
||||||
Bucket: bucket,
|
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 PO file(s) from S3`);
|
console.log(`Deleted ${keys.length} stale PO file(s) from S3`);
|
||||||
}
|
}
|
||||||
|
|
||||||
export const handler = async (): Promise<void> => {
|
export const handler = async (): Promise<void> => {
|
||||||
|
|
@ -169,49 +175,65 @@ export const handler = async (): Promise<void> => {
|
||||||
const kbId = process.env.KNOWLEDGE_BASE_ID!;
|
const kbId = process.env.KNOWLEDGE_BASE_ID!;
|
||||||
const dsId = process.env.DATA_SOURCE_ID!;
|
const dsId = process.env.DATA_SOURCE_ID!;
|
||||||
|
|
||||||
console.log('Clearing stale PO files from S3...');
|
console.log('Listing existing PO files in S3...');
|
||||||
await clearOldFiles(bucket);
|
const existingKeys = await listExistingKeys(bucket);
|
||||||
|
console.log(`Found ${existingKeys.size} existing file(s) in S3`);
|
||||||
|
|
||||||
console.log('Scanning purchase-orders table...');
|
console.log('Scanning purchase-orders table...');
|
||||||
const pos = await scanAllPOs();
|
const pos = await scanAllPOs();
|
||||||
console.log(`Found ${pos.length} purchase order(s)`);
|
console.log(`Found ${pos.length} purchase order(s)`);
|
||||||
|
|
||||||
|
const writtenKeys = new Set<string>();
|
||||||
let synced = 0;
|
let synced = 0;
|
||||||
let failed = 0;
|
let failed = 0;
|
||||||
|
|
||||||
for (const po of pos) {
|
// Upload in batches of 25 concurrent requests
|
||||||
const poNumber = po.po_number;
|
const CONCURRENCY = 25;
|
||||||
if (!poNumber) {
|
for (let i = 0; i < pos.length; i += CONCURRENCY) {
|
||||||
console.warn('Skipping record with no po_number');
|
const batch = pos.slice(i, i + CONCURRENCY);
|
||||||
failed++;
|
const results = await Promise.allSettled(
|
||||||
continue;
|
batch.map(async (po) => {
|
||||||
}
|
const poNumber = po.po_number;
|
||||||
|
if (!poNumber) {
|
||||||
|
console.warn('Skipping record with no po_number');
|
||||||
|
throw new Error('no po_number');
|
||||||
|
}
|
||||||
|
|
||||||
try {
|
const md = poToMarkdown(po);
|
||||||
const md = poToMarkdown(po);
|
const safeKey = poNumber.replace(/[^a-zA-Z0-9._-]/g, '_');
|
||||||
// Sanitize PO number for use as S3 key
|
const key = `${S3_PREFIX}${safeKey}.md`;
|
||||||
const safeKey = poNumber.replace(/[^a-zA-Z0-9._-]/g, '_');
|
|
||||||
|
|
||||||
await s3.send(
|
await s3.send(
|
||||||
new PutObjectCommand({
|
new PutObjectCommand({
|
||||||
Bucket: bucket,
|
Bucket: bucket,
|
||||||
Key: `${S3_PREFIX}${safeKey}.md`,
|
Key: key,
|
||||||
Body: md,
|
Body: md,
|
||||||
ContentType: 'text/markdown',
|
ContentType: 'text/markdown',
|
||||||
Metadata: { 'po-number': poNumber },
|
Metadata: { 'po-number': poNumber },
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
|
|
||||||
console.log(`Synced: PO ${poNumber}`);
|
return key;
|
||||||
synced++;
|
}),
|
||||||
} catch (err) {
|
);
|
||||||
console.error(`Failed to sync PO ${poNumber}:`, err);
|
|
||||||
failed++;
|
for (const r of results) {
|
||||||
|
if (r.status === 'fulfilled') {
|
||||||
|
writtenKeys.add(r.value);
|
||||||
|
synced++;
|
||||||
|
} else {
|
||||||
|
console.error(`Failed to sync PO:`, r.reason);
|
||||||
|
failed++;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
console.log(`Upload complete. ${synced} synced, ${failed} 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...');
|
console.log('Triggering Bedrock KB ingestion job...');
|
||||||
const ingestionRes = await bedrockAgent.send(
|
const ingestionRes = await bedrockAgent.send(
|
||||||
new StartIngestionJobCommand({
|
new StartIngestionJobCommand({
|
||||||
|
|
|
||||||
|
|
@ -33,7 +33,7 @@ export class PoSyncConstruct extends Construct {
|
||||||
entry: path.join(__dirname, '../../lambda/po-sync/index.ts'),
|
entry: path.join(__dirname, '../../lambda/po-sync/index.ts'),
|
||||||
runtime: lambda.Runtime.NODEJS_22_X,
|
runtime: lambda.Runtime.NODEJS_22_X,
|
||||||
memorySize: 512,
|
memorySize: 512,
|
||||||
timeout: cdk.Duration.minutes(5),
|
timeout: cdk.Duration.minutes(15),
|
||||||
environment: {
|
environment: {
|
||||||
PO_TABLE: poTable.tableName,
|
PO_TABLE: poTable.tableName,
|
||||||
KB_BUCKET_NAME: props.kbDocsBucket.bucketName,
|
KB_BUCKET_NAME: props.kbDocsBucket.bucketName,
|
||||||
|
|
|
||||||
Reference in a new issue