feat(catalog): add possibility to stream entity pages
adds a new method `streamEntityPages` that can be used to stream entity pages from catalog instead streaming entities one by one. this can be more efficient load wise but not usable for every use case. Signed-off-by: Hellgren Heikki <heikki.hellgren@op.fi>
This commit is contained in:
@@ -4,7 +4,7 @@
|
||||
'@backstage/plugin-catalog-node': minor
|
||||
---
|
||||
|
||||
Introduced new `streamEntities` async generator method for the catalog.
|
||||
Introduced new `streamEntities` and `streamEntityPages` async generator methods for the catalog.
|
||||
|
||||
Catalog API and Catalog Service now includes a `streamEntities` method that allows for streaming entities from the catalog.
|
||||
This method is designed to handle large datasets efficiently by processing entities in a stream rather than loading them
|
||||
@@ -19,3 +19,17 @@ for await (const entity of stream) {
|
||||
// Handle entity
|
||||
}
|
||||
```
|
||||
|
||||
Additionally, a `streamEntityPages` method is available that streams entities in pages, allowing for batch processing of entities.
|
||||
This method is more efficient than `streamEntities` when you can process entities in chunks.
|
||||
Example usage:
|
||||
|
||||
```ts
|
||||
const pageStream = catalogClient.streamEntityPages(
|
||||
{ batchSize: 100 },
|
||||
{ token },
|
||||
);
|
||||
for await (const page of pageStream) {
|
||||
// Handle page of entities
|
||||
}
|
||||
```
|
||||
|
||||
@@ -71,14 +71,9 @@ export class InMemoryCatalogClient implements CatalogApi {
|
||||
// (undocumented)
|
||||
removeLocationById(_id: string): Promise<void>;
|
||||
// (undocumented)
|
||||
streamEntities<
|
||||
T extends StreamEntitiesRequest,
|
||||
R = T extends {
|
||||
mode: 'array';
|
||||
} & StreamEntitiesRequest
|
||||
? Entity[]
|
||||
: Entity,
|
||||
>(request?: T): AsyncIterable<R>;
|
||||
streamEntities(request?: StreamEntitiesRequest): AsyncIterable<Entity>;
|
||||
// (undocumented)
|
||||
streamEntityPages(request?: StreamEntitiesRequest): AsyncIterable<Entity[]>;
|
||||
// (undocumented)
|
||||
validateEntity(
|
||||
_entity: Entity,
|
||||
|
||||
@@ -91,7 +91,11 @@ export interface CatalogApi {
|
||||
streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]>;
|
||||
): AsyncIterable<Entity>;
|
||||
streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
validateEntity(
|
||||
entity: Entity,
|
||||
locationRef: string,
|
||||
@@ -174,14 +178,14 @@ export class CatalogClient implements CatalogApi {
|
||||
id: string,
|
||||
options?: CatalogRequestOptions,
|
||||
): Promise<void>;
|
||||
streamEntities<
|
||||
T extends StreamEntitiesRequest,
|
||||
R = T extends {
|
||||
mode: 'array';
|
||||
} & StreamEntitiesRequest
|
||||
? Entity[]
|
||||
: Entity,
|
||||
>(request?: T, options?: CatalogRequestOptions): AsyncIterable<R>;
|
||||
streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity>;
|
||||
streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
validateEntity(
|
||||
entity: Entity,
|
||||
locationRef: string,
|
||||
@@ -335,7 +339,6 @@ export type StreamEntitiesRequest = Omit<
|
||||
'limit' | 'offset'
|
||||
> & {
|
||||
batchSize?: number;
|
||||
mode?: 'single' | 'array';
|
||||
};
|
||||
|
||||
// @public
|
||||
|
||||
@@ -579,7 +579,7 @@ describe('CatalogClient', () => {
|
||||
);
|
||||
});
|
||||
|
||||
it('should stream entities one by one', async () => {
|
||||
it('should stream entities', async () => {
|
||||
const stream = client.streamEntities({}, { token });
|
||||
const results: Entity[] = [];
|
||||
for await (const entity of stream) {
|
||||
@@ -588,13 +588,13 @@ describe('CatalogClient', () => {
|
||||
expect(results).toEqual(defaultResponse.items);
|
||||
});
|
||||
|
||||
it('should stream entities in batches', async () => {
|
||||
const stream = client.streamEntities({ mode: 'array' }, { token });
|
||||
const results: Entity[] = [];
|
||||
for await (const entityBatch of stream) {
|
||||
results.push(...entityBatch);
|
||||
it('should stream entity pages', async () => {
|
||||
const stream = client.streamEntityPages({}, { token });
|
||||
const results: Entity[][] = [];
|
||||
for await (const entityPage of stream) {
|
||||
results.push(entityPage);
|
||||
}
|
||||
expect(results).toEqual(defaultResponse.items);
|
||||
expect(results).toEqual([defaultResponse.items, []]);
|
||||
});
|
||||
|
||||
it('should handle errors', async () => {
|
||||
|
||||
@@ -463,26 +463,34 @@ export class CatalogClient implements CatalogApi {
|
||||
/**
|
||||
* {@inheritdoc CatalogApi.streamEntities}
|
||||
*/
|
||||
async *streamEntities<
|
||||
T extends StreamEntitiesRequest,
|
||||
R = T extends { mode: 'array' } & StreamEntitiesRequest ? Entity[] : Entity,
|
||||
>(request?: T, options?: CatalogRequestOptions): AsyncIterable<R> {
|
||||
async *streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity> {
|
||||
const pages = this.streamEntityPages(request, options);
|
||||
for await (const page of pages) {
|
||||
for (const entity of page) {
|
||||
yield entity;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc CatalogApi.streamEntityPages}
|
||||
*/
|
||||
async *streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity[]> {
|
||||
let cursor: string | undefined = undefined;
|
||||
const limit = request?.batchSize ?? DEFAULT_STREAM_ENTITIES_LIMIT;
|
||||
const mode = request?.mode ?? 'single';
|
||||
do {
|
||||
const res = await this.queryEntities(
|
||||
cursor ? { ...request, cursor, limit } : { ...request, limit },
|
||||
options,
|
||||
);
|
||||
|
||||
if (mode === 'single') {
|
||||
for (const entity of res.items) {
|
||||
yield entity as R;
|
||||
}
|
||||
} else {
|
||||
yield res.items as R;
|
||||
}
|
||||
yield res.items;
|
||||
|
||||
cursor = res.pageInfo.nextCursor;
|
||||
} while (cursor);
|
||||
|
||||
@@ -123,7 +123,7 @@ describe('InMemoryCatalogClient', () => {
|
||||
});
|
||||
});
|
||||
|
||||
it('streamEntities single', async () => {
|
||||
it('streamEntities', async () => {
|
||||
const client = new InMemoryCatalogClient({ entities });
|
||||
const stream = client.streamEntities();
|
||||
const results: Entity[] = [];
|
||||
@@ -133,14 +133,14 @@ describe('InMemoryCatalogClient', () => {
|
||||
expect(results).toEqual(entities);
|
||||
});
|
||||
|
||||
it('streamEntities batch', async () => {
|
||||
it('streamEntityPages', async () => {
|
||||
const client = new InMemoryCatalogClient({ entities });
|
||||
const stream = client.streamEntities({ mode: 'array' });
|
||||
const results: Entity[] = [];
|
||||
for await (const entity of stream) {
|
||||
results.push(...entity);
|
||||
const stream = client.streamEntityPages();
|
||||
const results: Entity[][] = [];
|
||||
for await (const page of stream) {
|
||||
results.push(page);
|
||||
}
|
||||
expect(results).toEqual(entities);
|
||||
expect(results).toEqual([entities]);
|
||||
});
|
||||
|
||||
it('getEntityAncestors', async () => {
|
||||
|
||||
@@ -280,24 +280,27 @@ export class InMemoryCatalogClient implements CatalogApi {
|
||||
throw new NotImplementedError('Method not implemented.');
|
||||
}
|
||||
|
||||
async *streamEntities<
|
||||
T extends StreamEntitiesRequest,
|
||||
R = T extends { mode: 'array' } & StreamEntitiesRequest ? Entity[] : Entity,
|
||||
>(request?: T): AsyncIterable<R> {
|
||||
async *streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
): AsyncIterable<Entity> {
|
||||
const pages = this.streamEntityPages(request);
|
||||
for await (const page of pages) {
|
||||
for (const entity of page) {
|
||||
yield entity;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async *streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
): AsyncIterable<Entity[]> {
|
||||
let cursor: string | undefined = undefined;
|
||||
const mode = request?.mode ?? 'single';
|
||||
do {
|
||||
const res = await this.queryEntities(
|
||||
cursor ? { ...request, cursor } : request,
|
||||
);
|
||||
|
||||
if (mode === 'single') {
|
||||
for (const entity of res.items) {
|
||||
yield entity as R;
|
||||
}
|
||||
} else {
|
||||
yield res.items as R;
|
||||
}
|
||||
yield res.items;
|
||||
|
||||
cursor = res.pageInfo.nextCursor;
|
||||
} while (cursor);
|
||||
|
||||
@@ -478,10 +478,6 @@ export type StreamEntitiesRequest = Omit<
|
||||
* The number of entities to fetch in each batch. Defaults to 500.
|
||||
*/
|
||||
batchSize?: number;
|
||||
/**
|
||||
* The mode in which entities should be yielded. Defaults to 'single'.
|
||||
*/
|
||||
mode?: 'single' | 'array';
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -724,5 +720,19 @@ export interface CatalogApi {
|
||||
streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]>;
|
||||
): AsyncIterable<Entity>;
|
||||
|
||||
/**
|
||||
* Asynchronously streams entity pages from the catalog. Uses `queryEntities`
|
||||
* to fetch entities in batches, and yields them one page at a time.
|
||||
*
|
||||
* @public
|
||||
*
|
||||
* @param request - Request parameters
|
||||
* @param options - Additional options
|
||||
*/
|
||||
streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
}
|
||||
|
||||
@@ -111,7 +111,12 @@ export interface CatalogServiceMock extends CatalogService, CatalogApi {
|
||||
streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogServiceRequestOptions | CatalogRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]>;
|
||||
): AsyncIterable<Entity>;
|
||||
// (undocumented)
|
||||
streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogServiceRequestOptions | CatalogRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
// (undocumented)
|
||||
validateEntity(
|
||||
entity: Entity,
|
||||
|
||||
@@ -199,7 +199,12 @@ export interface CatalogService {
|
||||
streamEntities(
|
||||
request: StreamEntitiesRequest | undefined,
|
||||
options: CatalogServiceRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]>;
|
||||
): AsyncIterable<Entity>;
|
||||
// (undocumented)
|
||||
streamEntityPages(
|
||||
request: StreamEntitiesRequest | undefined,
|
||||
options: CatalogServiceRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
// (undocumented)
|
||||
validateEntity(
|
||||
entity: Entity,
|
||||
|
||||
@@ -146,7 +146,12 @@ export interface CatalogService {
|
||||
streamEntities(
|
||||
request: StreamEntitiesRequest | undefined,
|
||||
options: CatalogServiceRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]>;
|
||||
): AsyncIterable<Entity>;
|
||||
|
||||
streamEntityPages(
|
||||
request: StreamEntitiesRequest | undefined,
|
||||
options: CatalogServiceRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
}
|
||||
|
||||
class DefaultCatalogService implements CatalogService {
|
||||
@@ -329,13 +334,23 @@ class DefaultCatalogService implements CatalogService {
|
||||
async *streamEntities(
|
||||
request: StreamEntitiesRequest | undefined,
|
||||
options: CatalogServiceRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]> {
|
||||
): AsyncIterable<Entity> {
|
||||
yield* this.#catalogApi.streamEntities(
|
||||
request,
|
||||
await this.#getOptions(options),
|
||||
);
|
||||
}
|
||||
|
||||
async *streamEntityPages(
|
||||
request: StreamEntitiesRequest | undefined,
|
||||
options: CatalogServiceRequestOptions,
|
||||
): AsyncIterable<Entity[]> {
|
||||
yield* this.#catalogApi.streamEntityPages(
|
||||
request,
|
||||
await this.#getOptions(options),
|
||||
);
|
||||
}
|
||||
|
||||
async #getOptions(
|
||||
options: CatalogServiceRequestOptions,
|
||||
): Promise<CatalogRequestOptions> {
|
||||
|
||||
@@ -104,5 +104,6 @@ export namespace catalogServiceMock {
|
||||
validateEntity: jest.fn(),
|
||||
analyzeLocation: jest.fn(),
|
||||
streamEntities: jest.fn(),
|
||||
streamEntityPages: jest.fn(),
|
||||
}));
|
||||
}
|
||||
|
||||
@@ -140,5 +140,10 @@ export interface CatalogServiceMock extends CatalogService, CatalogApi {
|
||||
streamEntities(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogServiceRequestOptions | CatalogRequestOptions,
|
||||
): AsyncIterable<Entity | Entity[]>;
|
||||
): AsyncIterable<Entity>;
|
||||
|
||||
streamEntityPages(
|
||||
request?: StreamEntitiesRequest,
|
||||
options?: CatalogServiceRequestOptions | CatalogRequestOptions,
|
||||
): AsyncIterable<Entity[]>;
|
||||
}
|
||||
|
||||
@@ -103,5 +103,6 @@ export namespace catalogApiMock {
|
||||
validateEntity: jest.fn(),
|
||||
analyzeLocation: jest.fn(),
|
||||
streamEntities: jest.fn(),
|
||||
streamEntityPages: jest.fn(),
|
||||
}));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user