diff --git a/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts b/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts index 56c6ac231a..84234d85a2 100644 --- a/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts +++ b/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts @@ -73,38 +73,6 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { return items.map(i => i.entity); } - private async addOrUpdateEntity( - entity: Entity, - tx: Transaction, - locationId?: string, - ): Promise { - // Find a matching (by uid, or by compound name, depending on the given - // entity) existing entity, to know whether to update or add - const existing = entity.metadata.uid - ? await this.database.entityByUid(tx, entity.metadata.uid) - : await this.database.entityByName(tx, getEntityName(entity)); - - // If it's an update, run the algorithm for annotation merging, updating - // etag/generation, etc. - let response: DbEntityResponse; - if (existing) { - const updated = generateUpdatedEntity(existing.entity, entity); - response = await this.database.updateEntity( - tx, - { locationId, entity: updated }, - existing.entity.metadata.etag, - existing.entity.metadata.generation, - ); - } else { - const added = await this.database.addEntities(tx, [ - { locationId, entity }, - ]); - response = added[0]; - } - - return response.entity; - } - async removeEntityByUid(uid: string): Promise { return await this.database.transaction(async tx => { const entityResponse = await this.database.entityByUid(tx, uid); @@ -135,13 +103,6 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { }); } - /** - * Writes a number of entities efficiently to storage. - * - * @param entities Some entities - * @param options.locationId The location that they all belong to - * @param options.tx A database transaction to execute the queries in - */ async batchAddOrUpdateEntities( requests: EntityUpsertRequest[], options?: { @@ -150,94 +111,33 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { outputEntities?: boolean; }, ): Promise { - const locationId = options?.locationId; - - // Group the entities by unique kind+namespace combinations + // Group the entities by unique kind+namespace combinations. The reason for + // this is that the change detection and merging logic requires finding + // pre-existing versions of the entities in the database. Those queries are + // easier and faster to make if every batch revolves around a single kind- + // namespace pair. const entitiesByKindAndNamespace = groupBy(requests, ({ entity }) => { const name = getEntityName(entity); return `${name.kind}:${name.namespace}`.toLowerCase(); }); + // Bound the number of concurrent batches. We want a bit of concurrency for + // performance reasons, but not so much that we starve the connection pool + // or start thrashing. const limiter = limiterFactory(BATCH_CONCURRENCY); - const tasks: Promise[] = []; + const tasks = new Array>(); for (const groupRequests of Object.values(entitiesByKindAndNamespace)) { - const { kind, namespace } = getEntityName(groupRequests[0].entity); - // Go through the new entities in reasonable chunk sizes (sometimes, // sources produce tens of thousands of entities, and those are too large // batch sizes to reasonably send to the database) for (const batch of chunk(groupRequests, BATCH_SIZE)) { tasks.push( limiter(async () => { - const first = serializeEntityRef(batch[0].entity); - const last = serializeEntityRef(batch[batch.length - 1].entity); - - this.logger.debug( - `Considering batch ${first}-${last} (${batch.length} entries)`, - ); - // Retry the batch write a few times to deal with contention - const context = { - kind, - namespace, - locationId, - }; - for (let attempt = 1; ; ++attempt) { try { - return await this.database.transaction(async tx => { - let modifiedEntityIds = new Array(); - const { toAdd, toUpdate, toIgnore } = await this.analyzeBatch( - batch, - context, - tx, - ); - - if (toAdd.length) { - modifiedEntityIds.push( - ...(await this.batchAdd(toAdd, context, tx)), - ); - } - if (toUpdate.length) { - modifiedEntityIds.push( - ...(await this.batchUpdate(toUpdate, context, tx)), - ); - } - - // TODO(Rugvip): We currently always update relations, but we - // likely want to figure out a way to avoid that - for (const { entity, relations } of toIgnore) { - const entityId = entity.metadata.uid; - if (entityId) { - await this.setRelations(entityId, relations, tx); - modifiedEntityIds.push({ entityId }); - } - } - - if (options?.outputEntities) { - const writtenEntities = await this.database.entities( - tx, - EntityFilters.ofMatchers({ - 'metadata.uid': modifiedEntityIds.map(e => e.entityId), - }), - ); - modifiedEntityIds = writtenEntities.map(e => ({ - entityId: e.entity.metadata.uid!, - entity: e.entity, - })); - } - - if (options?.dryRun) { - // If this is only a dry run, cancel the database transaction even if it was successful. - await tx.rollback(); - this.logger.debug( - `Performed successful dry run of adding entities`, - ); - } - - return modifiedEntityIds; - }); + return this.batchAddOrUpdateEntitiesSingleBatch(batch, options); } catch (e) { if (e instanceof ConflictError && attempt < BATCH_ATTEMPTS) { this.logger.warn( @@ -253,8 +153,83 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { } } - const entityUpserts = await Promise.all(tasks); - return entityUpserts.flat(); + const responses = await Promise.all(tasks); + return responses.flat(); + } + + // Defines the actual logic of running a single batch. All of these share a + // common kind and namespace. + private async batchAddOrUpdateEntitiesSingleBatch( + batch: EntityUpsertRequest[], + options?: { + locationId?: string; + dryRun?: boolean; + outputEntities?: boolean; + }, + ) { + const { kind, namespace } = getEntityName(batch[0].entity); + const context = { + kind, + namespace, + locationId: options?.locationId, + }; + + this.logger.debug( + `Considering batch ${serializeEntityRef( + batch[0].entity, + )}-${serializeEntityRef(batch[batch.length - 1].entity)} (${ + batch.length + } entries)`, + ); + + return this.database.transaction(async tx => { + const { toAdd, toUpdate, toIgnore } = await this.analyzeBatch( + batch, + context, + tx, + ); + + let responses = new Array(); + if (toAdd.length) { + const items = await this.batchAdd(toAdd, context, tx); + responses.push(...items); + } + if (toUpdate.length) { + const items = await this.batchUpdate(toUpdate, context, tx); + responses.push(...items); + } + for (const { entity, relations } of toIgnore) { + // TODO(Rugvip): We currently always update relations, but we + // likely want to figure out a way to avoid that + const entityId = entity.metadata.uid; + if (entityId) { + await this.setRelations(entityId, relations, tx); + responses.push({ entityId }); + } + } + + if (options?.outputEntities) { + const writtenEntities = await this.database.entities( + tx, + EntityFilters.ofMatchers({ + 'metadata.uid': responses.map(e => e.entityId), + }), + ); + responses = writtenEntities.map(e => ({ + entityId: e.entity.metadata.uid!, + entity: e.entity, + })); + } + + // If this is only a dry run, cancel the database transaction even if it + // was successful. + if (options?.dryRun) { + await tx.rollback(); + this.logger.debug(`Performed successful dry run of adding entities`); + } + + return responses; + }); } // Set the relations originating from an entity using the DB layer @@ -281,6 +256,8 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { }> { const markTimestamp = process.hrtime(); + // Here we make use of the fact that all of the entities share kind and + // namespace within a batch const names = requests.map(({ entity }) => entity.metadata.name); const oldEntities = await this.database.entities( tx, @@ -347,11 +324,11 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { requests.map(({ entity }) => ({ locationId, entity })), ); - const entityIds = res.map(({ entity }) => ({ + const responses = res.map(({ entity }) => ({ entityId: entity.metadata.uid!, })); - for (const [index, { entityId }] of entityIds.entries()) { + for (const [index, { entityId }] of responses.entries()) { await this.setRelations(entityId, requests[index].relations, tx); } @@ -359,7 +336,7 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { `Added ${requests.length} entities in ${durationText(markTimestamp)}`, ); - return entityIds; + return responses; } // Efficiently updates the given entities into storage, under the assumption @@ -370,12 +347,13 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { tx: Transaction, ): Promise { const markTimestamp = process.hrtime(); - const responseIds: EntityUpsertResponse[] = []; + const responses: EntityUpsertResponse[] = []; + // TODO(freben): Still not batched for (const entity of requests) { const res = await this.addOrUpdateEntity(entity.entity, tx, locationId); const entityId = res.metadata.uid!; - responseIds.push({ entityId }); + responses.push({ entityId }); await this.setRelations(entityId, entity.relations, tx); } @@ -383,6 +361,39 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { `Updated ${requests.length} entities in ${durationText(markTimestamp)}`, ); - return responseIds; + return responses; + } + + // TODO(freben): Incorporate this into batchUpdate which is the only caller + private async addOrUpdateEntity( + entity: Entity, + tx: Transaction, + locationId?: string, + ): Promise { + // Find a matching (by uid, or by compound name, depending on the given + // entity) existing entity, to know whether to update or add + const existing = entity.metadata.uid + ? await this.database.entityByUid(tx, entity.metadata.uid) + : await this.database.entityByName(tx, getEntityName(entity)); + + // If it's an update, run the algorithm for annotation merging, updating + // etag/generation, etc. + let response: DbEntityResponse; + if (existing) { + const updated = generateUpdatedEntity(existing.entity, entity); + response = await this.database.updateEntity( + tx, + { locationId, entity: updated }, + existing.entity.metadata.etag, + existing.entity.metadata.generation, + ); + } else { + const added = await this.database.addEntities(tx, [ + { locationId, entity }, + ]); + response = added[0]; + } + + return response.entity; } } diff --git a/plugins/catalog-backend/src/catalog/types.ts b/plugins/catalog-backend/src/catalog/types.ts index a3b2e0cb96..b012ee3606 100644 --- a/plugins/catalog-backend/src/catalog/types.ts +++ b/plugins/catalog-backend/src/catalog/types.ts @@ -32,17 +32,30 @@ export type EntityUpsertResponse = { }; export type EntitiesCatalog = { + /** + * Fetch entities. + * + * @param filter A filter to apply when reading + */ entities(filter?: EntityFilter): Promise; + + /** + * Removes a single entity. + * + * @param uid The metadata.uid of the entity + */ removeEntityByUid(uid: string): Promise; /** * Writes a number of entities efficiently to storage. * - * @param entities Some entities - * @param locationId The location that they all belong to + * @param requests The entities and their relations + * @param options.locationId The location that they all belong to (default none) + * @param options.dryRun Whether to throw away the results (default false) + * @param options.outputEntities Whether to return the resulting entities (default false) */ batchAddOrUpdateEntities( - entities: EntityUpsertRequest[], + requests: EntityUpsertRequest[], options?: { locationId?: string; dryRun?: boolean;