catalog-backend: rearrange and comment a bit
This commit is contained in:
@@ -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<Entity> {
|
||||
// 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<void> {
|
||||
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<EntityUpsertResponse[]> {
|
||||
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<EntityUpsertResponse[]>[] = [];
|
||||
const tasks = new Array<Promise<EntityUpsertResponse[]>>();
|
||||
|
||||
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<EntityUpsertResponse>();
|
||||
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<EntityUpsertResponse>();
|
||||
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<EntityUpsertResponse[]> {
|
||||
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<Entity> {
|
||||
// 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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,17 +32,30 @@ export type EntityUpsertResponse = {
|
||||
};
|
||||
|
||||
export type EntitiesCatalog = {
|
||||
/**
|
||||
* Fetch entities.
|
||||
*
|
||||
* @param filter A filter to apply when reading
|
||||
*/
|
||||
entities(filter?: EntityFilter): Promise<Entity[]>;
|
||||
|
||||
/**
|
||||
* Removes a single entity.
|
||||
*
|
||||
* @param uid The metadata.uid of the entity
|
||||
*/
|
||||
removeEntityByUid(uid: string): Promise<void>;
|
||||
|
||||
/**
|
||||
* 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;
|
||||
|
||||
Reference in New Issue
Block a user