Added support for multiple consumed topics

This commit is contained in:
Nir Gazit
2021-01-21 12:32:02 +02:00
parent d87367a6b3
commit 7112565f45
7 changed files with 158 additions and 56 deletions
@@ -15,7 +15,7 @@
*/
import { Table, TableColumn } from '@backstage/core';
import { Box, Typography } from '@material-ui/core';
import { Box, Grid, Typography } from '@material-ui/core';
import RetryIcon from '@material-ui/icons/Replay';
import React from 'react';
import { useConsumerGroupsOffsetsForEntity } from './useConsumerGroupsOffsetsForEntity';
@@ -62,6 +62,7 @@ const generatedColumns: TableColumn[] = [
type Props = {
loading: boolean;
retry: () => void;
clusterId: string;
consumerGroup: string;
topics?: TopicPartitionInfo[];
};
@@ -69,6 +70,7 @@ type Props = {
export const ConsumerGroupOffsets = ({
loading,
topics,
clusterId,
consumerGroup,
retry,
}: Props) => {
@@ -87,7 +89,7 @@ export const ConsumerGroupOffsets = ({
title={
<Box display="flex" alignItems="center">
<Typography variant="h6">
Consumed Topics for {consumerGroup}
Consumed Topics for {consumerGroup} ({clusterId})
</Typography>
</Box>
}
@@ -98,6 +100,15 @@ export const ConsumerGroupOffsets = ({
export const KafkaTopicsForConsumer = () => {
const [tableProps, { retry }] = useConsumerGroupsOffsetsForEntity();
return <ConsumerGroupOffsets {...tableProps} retry={retry} />;
return (
<Grid>
{tableProps.consumerGroupsTopics?.map(consumerGroup => (
<ConsumerGroupOffsets
{...consumerGroup}
loading={tableProps.loading}
retry={retry}
/>
))}
</Grid>
);
};
@@ -49,10 +49,66 @@ describe('useConsumerGroupOffsets', () => {
},
};
const { result } = subject();
expect(result.current).toStrictEqual({
clusterId: 'prod',
consumerGroup: 'consumer',
});
expect(result.current).toStrictEqual([
{
clusterId: 'prod',
consumerGroup: 'consumer',
},
]);
});
it('returns correct cluster and consumer group for multiple consumers', async () => {
entity = {
apiVersion: 'v1',
kind: 'Component',
metadata: {
name: 'test',
annotations: {
'kafka.apache.org/consumer-groups':
'prod/consumer,dev/another-consumer',
},
},
spec: {
owner: 'guest',
type: 'Website',
lifecycle: 'development',
},
};
const { result } = subject();
expect(result.current).toStrictEqual([
{ clusterId: 'prod', consumerGroup: 'consumer' },
{
clusterId: 'dev',
consumerGroup: 'another-consumer',
},
]);
});
it('returns correct cluster and consumer group for annotation with extra spaces', async () => {
entity = {
apiVersion: 'v1',
kind: 'Component',
metadata: {
name: 'test',
annotations: {
'kafka.apache.org/consumer-groups':
' prod/consumer , dev/another-consumer ',
},
},
spec: {
owner: 'guest',
type: 'Website',
lifecycle: 'development',
},
};
const { result } = subject();
expect(result.current).toStrictEqual([
{ clusterId: 'prod', consumerGroup: 'consumer' },
{
clusterId: 'dev',
consumerGroup: 'another-consumer',
},
]);
});
it('fails on missing cluster', async () => {
@@ -62,7 +118,7 @@ describe('useConsumerGroupOffsets', () => {
metadata: {
name: 'test',
annotations: {
'kafka.apache.org/consumer-groups': 'consumer',
'kafka.apache.org/consumer-groups': 'dev/another,consumer',
},
},
spec: {
@@ -75,7 +131,7 @@ describe('useConsumerGroupOffsets', () => {
expect(() => result.current).toThrowError();
expect(result.error).toStrictEqual(
new Error(
`Failed to parse kafka consumer group annotation: got "consumer"`,
`Failed to parse kafka consumer group annotation: got "dev/another,consumer"`,
),
);
});
@@ -15,19 +15,29 @@
*/
import { useEntity } from '@backstage/plugin-catalog';
import { useMemo } from 'react';
import { KAFKA_CONSUMER_GROUP_ANNOTATION } from '../../constants';
export const useConsumerGroupsForEntity = () => {
const { entity } = useEntity();
const annotation =
entity.metadata.annotations?.[KAFKA_CONSUMER_GROUP_ANNOTATION] ?? '';
const [clusterId, consumerGroup] = annotation.split('/');
if (!clusterId || !consumerGroup) {
throw new Error(
`Failed to parse kafka consumer group annotation: got "${annotation}"`,
);
}
const consumerList = useMemo(() => {
return annotation.split(',').map(consumer => {
const [clusterId, consumerGroup] = consumer.split('/');
return { clusterId, consumerGroup };
if (!clusterId || !consumerGroup) {
throw new Error(
`Failed to parse kafka consumer group annotation: got "${annotation}"`,
);
}
return {
clusterId: clusterId.trim(),
consumerGroup: consumerGroup.trim(),
};
});
}, [annotation]);
return consumerList;
};
@@ -45,7 +45,7 @@ describe('useConsumerGroupOffsets', () => {
metadata: {
name: 'test',
annotations: {
'kafka.apache.org/consumer-groups': `cluster/${consumerGroupOffsets.consumerId}`,
'kafka.apache.org/consumer-groups': `prod/${consumerGroupOffsets.consumerId}`,
},
},
spec: {
@@ -74,16 +74,24 @@ describe('useConsumerGroupOffsets', () => {
renderHook(useConsumerGroupsOffsetsForEntity, { wrapper });
it('returns correct consumer group for annotation', async () => {
mockKafkaApi.getConsumerGroupOffsets.mockResolvedValue(
consumerGroupOffsets,
);
when(mockKafkaApi.getConsumerGroupOffsets)
.calledWith('cluster', consumerGroupOffsets.consumerId)
.calledWith('prod', consumerGroupOffsets.consumerId)
.mockResolvedValue(consumerGroupOffsets);
const { result, waitForNextUpdate } = subject();
await waitForNextUpdate();
const [tableProps] = result.current;
expect(tableProps.consumerGroup).toBe(consumerGroupOffsets.consumerId);
expect(tableProps.topics).toBe(consumerGroupOffsets.offsets);
expect(tableProps.consumerGroupsTopics).toStrictEqual([
{
clusterId: 'prod',
consumerGroup: consumerGroupOffsets.consumerId,
topics: consumerGroupOffsets.offsets,
},
]);
});
it('posts an error to the error api', async () => {
@@ -20,28 +20,35 @@ import { kafkaApiRef } from '../../api/types';
import { useConsumerGroupsForEntity } from './useConsumerGroupsForEntity';
export const useConsumerGroupsOffsetsForEntity = () => {
const { clusterId, consumerGroup } = useConsumerGroupsForEntity();
const consumers = useConsumerGroupsForEntity();
const api = useApi(kafkaApiRef);
const errorApi = useApi(errorApiRef);
const { loading, value: topics, retry } = useAsyncRetry(async () => {
const {
loading,
value: consumerGroupsTopics,
retry,
} = useAsyncRetry(async () => {
try {
const response = await api.getConsumerGroupOffsets(
clusterId,
consumerGroup,
return await Promise.all(
consumers.map(async ({ clusterId, consumerGroup }) => {
const response = await api.getConsumerGroupOffsets(
clusterId,
consumerGroup,
);
return { clusterId, consumerGroup, topics: response.offsets };
}),
);
return response.offsets;
} catch (e) {
errorApi.post(e);
throw e;
}
}, [api, errorApi, consumerGroup]);
}, [consumers, api, errorApi]);
return [
{
loading,
consumerGroup,
topics,
consumerGroupsTopics,
},
{
retry,