From 285d21640318741b241579e5ee070fc343ed0db9 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Tue, 19 Jan 2021 16:59:39 +0200 Subject: [PATCH 1/9] Added support for multiple kafka clusters --- plugins/kafka-backend/README.md | 6 +- plugins/kafka-backend/config.d.ts | 15 ++-- .../src/config/ClusterReader.test.ts | 45 ++++++++++ .../kafka-backend/src/config/ClusterReader.ts | 37 ++++++++ .../kafka-backend/src/service/router.test.ts | 89 ++++++++++++++++--- plugins/kafka-backend/src/service/router.ts | 37 +++++--- plugins/kafka-backend/src/types/types.ts | 23 +++++ plugins/kafka/README.md | 8 +- plugins/kafka/src/api/KafkaBackendClient.ts | 5 +- plugins/kafka/src/api/types.ts | 1 + .../useConsumerGroupsForEntity.test.tsx | 63 +++++++++---- .../useConsumerGroupsForEntity.ts | 11 ++- ...useConsumerGroupsOffsetsForEntity.test.tsx | 4 +- .../useConsumerGroupsOffsetsForEntity.ts | 7 +- 14 files changed, 295 insertions(+), 56 deletions(-) create mode 100644 plugins/kafka-backend/src/config/ClusterReader.test.ts create mode 100644 plugins/kafka-backend/src/config/ClusterReader.ts create mode 100644 plugins/kafka-backend/src/types/types.ts diff --git a/plugins/kafka-backend/README.md b/plugins/kafka-backend/README.md index b2271a2911..13078ea816 100644 --- a/plugins/kafka-backend/README.md +++ b/plugins/kafka-backend/README.md @@ -23,6 +23,8 @@ Example: ```yaml kafka: clientId: backstage - brokers: - - localhost:9092 + clusters: + name: prod + brokers: + - localhost:9092 ``` diff --git a/plugins/kafka-backend/config.d.ts b/plugins/kafka-backend/config.d.ts index 0eeeb42197..122e76487a 100644 --- a/plugins/kafka-backend/config.d.ts +++ b/plugins/kafka-backend/config.d.ts @@ -16,11 +16,14 @@ export interface Config { kafka?: { clientId: string; - brokers: string[]; - ssl?: { - ca: string[]; - key: string; - cert: string; - }; + clusters: { + name: string; + brokers: string[]; + ssl?: { + ca: string[]; + key: string; + cert: string; + }; + }[]; }; } diff --git a/plugins/kafka-backend/src/config/ClusterReader.test.ts b/plugins/kafka-backend/src/config/ClusterReader.test.ts new file mode 100644 index 0000000000..0529684425 --- /dev/null +++ b/plugins/kafka-backend/src/config/ClusterReader.test.ts @@ -0,0 +1,45 @@ +/* + * Copyright 2020 Spotify AB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import '@backstage/backend-common'; +import { ConfigReader } from '@backstage/config'; +import { getClusterDetails } from './ClusterReader'; + +describe('getClusterDetails', () => { + it('empty clusters return empty cluster details', async () => { + const config = new ConfigReader({ clusters: [] }); + + const result = getClusterDetails(config.getConfigArray('clusters')); + + expect(result).toStrictEqual([]); + }); + + it('two clusters return two cluster details', async () => { + const config = new ConfigReader({ + clusters: [ + { name: 'cluster1', brokers: ['a', 'b'] }, + { name: 'cluster2', brokers: ['d'] }, + ], + }); + + const result = getClusterDetails(config.getConfigArray('clusters')); + + expect(result).toStrictEqual([ + { name: 'cluster1', brokers: ['a', 'b'] }, + { name: 'cluster2', brokers: ['d'] }, + ]); + }); +}); diff --git a/plugins/kafka-backend/src/config/ClusterReader.ts b/plugins/kafka-backend/src/config/ClusterReader.ts new file mode 100644 index 0000000000..a4e6a4e4fe --- /dev/null +++ b/plugins/kafka-backend/src/config/ClusterReader.ts @@ -0,0 +1,37 @@ +/* + * Copyright 2020 Spotify AB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { Config } from '@backstage/config'; +import { ConnectionOptions } from 'tls'; +import { ClusterDetails } from '../types/types'; + +export function getClusterDetails(config: Config[]): ClusterDetails[] { + return config.map(clusterConfig => { + const clusterDetails = { + name: clusterConfig.getString('name'), + brokers: clusterConfig.getStringArray('brokers'), + }; + const sslConfig = clusterConfig.getOptional('kafka.ssl'); + if (sslConfig) { + return { + ...clusterDetails, + ssl: sslConfig as ConnectionOptions, + }; + } + + return clusterDetails; + }); +} diff --git a/plugins/kafka-backend/src/service/router.test.ts b/plugins/kafka-backend/src/service/router.test.ts index 5e9de6cf2f..83eacc5b48 100644 --- a/plugins/kafka-backend/src/service/router.test.ts +++ b/plugins/kafka-backend/src/service/router.test.ts @@ -16,22 +16,34 @@ import request from 'supertest'; import express from 'express'; -import { makeRouter } from './router'; +import { makeRouter, ClusterApi } from './router'; import { getVoidLogger } from '@backstage/backend-common'; import { KafkaApi } from './KafkaApi'; import { when } from 'jest-when'; describe('router', () => { let app: express.Express; - let kafkaApi: jest.Mocked; + let apis: ClusterApi[]; + let devKafkaApi: jest.Mocked; + let prodKafkaApi: jest.Mocked; beforeAll(async () => { - kafkaApi = { + devKafkaApi = { fetchTopicOffsets: jest.fn(), fetchGroupOffsets: jest.fn(), }; - const router = makeRouter(getVoidLogger(), kafkaApi); + prodKafkaApi = { + fetchTopicOffsets: jest.fn(), + fetchGroupOffsets: jest.fn(), + }; + + apis = [ + { name: 'dev', api: devKafkaApi }, + { name: 'prod', api: prodKafkaApi }, + ]; + + const router = makeRouter(getVoidLogger(), apis); app = express().use(router); }); @@ -39,7 +51,7 @@ describe('router', () => { jest.resetAllMocks(); }); - describe('get /consumer/:consumerId/offsets', () => { + describe('get /:clusterId/consumer/:consumerId/offsets', () => { it('returns topic and group offsets', async () => { const topic1Offsets = [ { id: 1, offset: '500' }, @@ -60,15 +72,17 @@ describe('router', () => { partitions: [{ id: 1, offset: '456' }], }, ]; - when(kafkaApi.fetchTopicOffsets) + when(prodKafkaApi.fetchTopicOffsets) .calledWith('topic1') .mockResolvedValue(topic1Offsets); - when(kafkaApi.fetchTopicOffsets) + when(prodKafkaApi.fetchTopicOffsets) .calledWith('topic2') .mockResolvedValue(topic2Offsets); - kafkaApi.fetchGroupOffsets.mockResolvedValue(groupOffsets); + when(prodKafkaApi.fetchGroupOffsets) + .calledWith('hey') + .mockResolvedValue(groupOffsets); - const response = await request(app).get('/consumer/hey/offsets'); + const response = await request(app).get('/prod/consumer/hey/offsets'); expect(response.status).toEqual(200); expect(response.body.consumerId).toEqual('hey'); @@ -95,11 +109,64 @@ describe('router', () => { }); it('handles internal error correctly', async () => { - kafkaApi.fetchGroupOffsets.mockRejectedValue(Error('oh no')); + prodKafkaApi.fetchGroupOffsets.mockRejectedValue(Error('oh no')); - const response = await request(app).get('/consumer/hey/offsets'); + const response = await request(app).get('/prod/consumer/hey/offsets'); expect(response.status).toEqual(500); }); + + it('uses correct kafka cluster', async () => { + const topic1ProdOffsets = [{ id: 1, offset: '500' }]; + const topic1DevOffsets = [{ id: 1, offset: '1234' }]; + const groupProdOffsets = [ + { + topic: 'topic1', + partitions: [{ id: 1, offset: '100' }], + }, + ]; + const groupDevOffsets = [ + { + topic: 'topic1', + partitions: [{ id: 1, offset: '567' }], + }, + ]; + when(prodKafkaApi.fetchTopicOffsets) + .calledWith('topic1') + .mockResolvedValue(topic1ProdOffsets); + when(prodKafkaApi.fetchGroupOffsets) + .calledWith('hey') + .mockResolvedValue(groupProdOffsets); + when(devKafkaApi.fetchTopicOffsets) + .calledWith('topic1') + .mockResolvedValue(topic1DevOffsets); + when(devKafkaApi.fetchGroupOffsets) + .calledWith('hey') + .mockResolvedValue(groupDevOffsets); + + const prodResponse = await request(app).get('/prod/consumer/hey/offsets'); + const devResponse = await request(app).get('/dev/consumer/hey/offsets'); + + expect(prodResponse.status).toEqual(200); + expect(prodResponse.body.consumerId).toEqual('hey'); + expect(prodResponse.body.offsets).toIncludeSameMembers([ + { + topic: 'topic1', + partitionId: 1, + groupOffset: '100', + topicOffset: '500', + }, + ]); + expect(devResponse.status).toEqual(200); + expect(devResponse.body.consumerId).toEqual('hey'); + expect(devResponse.body.offsets).toIncludeSameMembers([ + { + topic: 'topic1', + partitionId: 1, + groupOffset: '567', + topicOffset: '1234', + }, + ]); + }); }); }); diff --git a/plugins/kafka-backend/src/service/router.ts b/plugins/kafka-backend/src/service/router.ts index 2bbc877fe2..d6e5ab08a7 100644 --- a/plugins/kafka-backend/src/service/router.ts +++ b/plugins/kafka-backend/src/service/router.ts @@ -20,31 +20,43 @@ import { Logger } from 'winston'; import { Config } from '@backstage/config'; import { KafkaApi, KafkaJsApiImpl } from './KafkaApi'; import _ from 'lodash'; -import { ConnectionOptions } from 'tls'; +import { getClusterDetails } from '../config/clusterReader'; export interface RouterOptions { logger: Logger; config: Config; } +export interface ClusterApi { + name: string; + api: KafkaApi; +} + export const makeRouter = ( logger: Logger, - kafkaApi: KafkaApi, + kafkaApis: ClusterApi[], ): express.Router => { const router = Router(); router.use(express.json()); - router.get('/consumer/:consumerId/offsets', async (req, res) => { + const kafkaApiByClusterName = _.keyBy(kafkaApis, item => item.name); + + router.get('/consumers/:clusterId/:consumerId/offsets', async (req, res) => { + const clusterId = req.params.clusterId; const consumerId = req.params.consumerId; - logger.debug(`Fetch consumer group ${consumerId} offsets`); + logger.info( + `Fetch consumer group ${consumerId} offsets from cluster ${clusterId}`, + ); - const groupOffsets = await kafkaApi.fetchGroupOffsets(consumerId); + const kafkaApi = kafkaApiByClusterName[clusterId]; + + const groupOffsets = await kafkaApi.api.fetchGroupOffsets(consumerId); const groupWithTopicOffsets = await Promise.all( groupOffsets.map(async ({ topic, partitions }) => { const topicOffsets = _.keyBy( - await kafkaApi.fetchTopicOffsets(topic), + await kafkaApi.api.fetchTopicOffsets(topic), partition => partition.id, ); @@ -70,12 +82,15 @@ export async function createRouter( logger.info('Initializing Kafka backend'); const clientId = options.config.getString('kafka.clientId'); - const brokers = options.config.getStringArray('kafka.brokers'); - const sslConfig = options.config.getOptional('kafka.ssl'); - const ssl = sslConfig ? (sslConfig as ConnectionOptions) : undefined; + const clusters = getClusterDetails( + options.config.getConfigArray('kafka.clusters'), + ); - const kafkaApi = new KafkaJsApiImpl({ clientId, brokers, logger, ssl }); + const kafkaApis = clusters.map(cluster => ({ + name: cluster.name, + api: new KafkaJsApiImpl({ clientId, logger, ...cluster }), + })); - return makeRouter(logger, kafkaApi); + return makeRouter(logger, kafkaApis); } diff --git a/plugins/kafka-backend/src/types/types.ts b/plugins/kafka-backend/src/types/types.ts new file mode 100644 index 0000000000..5f70998dd8 --- /dev/null +++ b/plugins/kafka-backend/src/types/types.ts @@ -0,0 +1,23 @@ +/* + * Copyright 2020 Spotify AB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { ConnectionOptions } from 'tls'; + +export interface ClusterDetails { + name: string; + brokers: string[]; + ssl?: ConnectionOptions; +} diff --git a/plugins/kafka/README.md b/plugins/kafka/README.md index 7d947ae41b..c41c328741 100644 --- a/plugins/kafka/README.md +++ b/plugins/kafka/README.md @@ -65,8 +65,10 @@ import { Router as KafkaRouter } from '@backstage/plugin-kafka'; ```yaml kafka: clientId: backstage - brokers: - - localhost:9092 + clusters: + - name: cluster-name + brokers: + - localhost:9092 ``` 6. Add `kafka.apache.org/consumer-groups` annotation to your services: @@ -77,7 +79,7 @@ kind: Component metadata: # ... annotations: - kafka.apache.org/consumer-groups: consumer-group-name + kafka.apache.org/consumer-groups: cluster-name/consumer-group-name spec: type: service ``` diff --git a/plugins/kafka/src/api/KafkaBackendClient.ts b/plugins/kafka/src/api/KafkaBackendClient.ts index 2112dd0f17..45b1616fe0 100644 --- a/plugins/kafka/src/api/KafkaBackendClient.ts +++ b/plugins/kafka/src/api/KafkaBackendClient.ts @@ -43,8 +43,11 @@ export class KafkaBackendClient implements KafkaApi { } async getConsumerGroupOffsets( + clusterId: string, consumerGroup: string, ): Promise { - return await this.getRequired(`/consumer/${consumerGroup}/offsets`); + return await this.getRequired( + `/consumers/${clusterId}/${consumerGroup}/offsets`, + ); } } diff --git a/plugins/kafka/src/api/types.ts b/plugins/kafka/src/api/types.ts index 9c2239156b..c4307d5c1f 100644 --- a/plugins/kafka/src/api/types.ts +++ b/plugins/kafka/src/api/types.ts @@ -34,6 +34,7 @@ export type ConsumerGroupOffsetsResponse = { export interface KafkaApi { getConsumerGroupOffsets( + clusterId: string, consumerGroup: string, ): Promise; } diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx index 44693d5bbc..483986d4ec 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx @@ -20,21 +20,7 @@ import { EntityContext } from '@backstage/plugin-catalog'; import { Entity } from '@backstage/catalog-model'; describe('useConsumerGroupOffsets', () => { - const entity: Entity = { - apiVersion: 'v1', - kind: 'Component', - metadata: { - name: 'test', - annotations: { - 'kafka.apache.org/consumer-groups': 'consumer', - }, - }, - spec: { - owner: 'guest', - type: 'Website', - lifecycle: 'development', - }, - }; + let entity: Entity; const wrapper = ({ children }: PropsWithChildren<{}>) => { return ( @@ -46,8 +32,51 @@ describe('useConsumerGroupOffsets', () => { const subject = () => renderHook(useConsumerGroupsForEntity, { wrapper }); - it('returns correct consumer group for annotation', async () => { + it('returns correct cluster and consumer group for annotation', async () => { + entity = { + apiVersion: 'v1', + kind: 'Component', + metadata: { + name: 'test', + annotations: { + 'kafka.apache.org/consumer-groups': 'prod/consumer', + }, + }, + spec: { + owner: 'guest', + type: 'Website', + lifecycle: 'development', + }, + }; const { result } = subject(); - expect(result.current).toBe('consumer'); + expect(result.current).toStrictEqual({ + clusterId: 'prod', + consumerGroup: 'consumer', + }); + }); + + it('fails on missing cluster', async () => { + entity = { + apiVersion: 'v1', + kind: 'Component', + metadata: { + name: 'test', + annotations: { + 'kafka.apache.org/consumer-groups': 'consumer', + }, + }, + spec: { + owner: 'guest', + type: 'Website', + lifecycle: 'development', + }, + }; + const { result } = subject(); + expect(() => result.current).toThrowError(); + expect(result.error).toStrictEqual( + new Error( + `Failed to parse kafka consumer group annotation: got "consumer"`, + ), + ); }); }); diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts index 5003c2f5b5..0afec42870 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts @@ -19,6 +19,15 @@ 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('/'); - return entity.metadata.annotations?.[KAFKA_CONSUMER_GROUP_ANNOTATION] ?? ''; + if (!clusterId || !consumerGroup) { + throw new Error( + `Failed to parse kafka consumer group annotation: got "${annotation}"`, + ); + } + + return { clusterId, consumerGroup }; }; diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx index 28f0ff7c28..16bf2f7665 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx @@ -45,7 +45,7 @@ describe('useConsumerGroupOffsets', () => { metadata: { name: 'test', annotations: { - 'kafka.apache.org/consumer-groups': consumerGroupOffsets.consumerId, + 'kafka.apache.org/consumer-groups': `cluster/${consumerGroupOffsets.consumerId}`, }, }, spec: { @@ -75,7 +75,7 @@ describe('useConsumerGroupOffsets', () => { it('returns correct consumer group for annotation', async () => { when(mockKafkaApi.getConsumerGroupOffsets) - .calledWith(consumerGroupOffsets.consumerId) + .calledWith('cluster', consumerGroupOffsets.consumerId) .mockResolvedValue(consumerGroupOffsets); const { result, waitForNextUpdate } = subject(); diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts index eb7e105989..4081b59796 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts @@ -20,13 +20,16 @@ import { kafkaApiRef } from '../../api/types'; import { useConsumerGroupsForEntity } from './useConsumerGroupsForEntity'; export const useConsumerGroupsOffsetsForEntity = () => { - const consumerGroup = useConsumerGroupsForEntity(); + const { clusterId, consumerGroup } = useConsumerGroupsForEntity(); const api = useApi(kafkaApiRef); const errorApi = useApi(errorApiRef); const { loading, value: topics, retry } = useAsyncRetry(async () => { try { - const response = await api.getConsumerGroupOffsets(consumerGroup); + const response = await api.getConsumerGroupOffsets( + clusterId, + consumerGroup, + ); return response.offsets; } catch (e) { errorApi.post(e); From 5c0fb3b8443781ff874def5911d07dad4567f3a4 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Tue, 19 Jan 2021 17:17:55 +0200 Subject: [PATCH 2/9] Fixed incorrect import --- plugins/kafka-backend/src/service/router.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/plugins/kafka-backend/src/service/router.ts b/plugins/kafka-backend/src/service/router.ts index d6e5ab08a7..465347ed8f 100644 --- a/plugins/kafka-backend/src/service/router.ts +++ b/plugins/kafka-backend/src/service/router.ts @@ -20,7 +20,7 @@ import { Logger } from 'winston'; import { Config } from '@backstage/config'; import { KafkaApi, KafkaJsApiImpl } from './KafkaApi'; import _ from 'lodash'; -import { getClusterDetails } from '../config/clusterReader'; +import { getClusterDetails } from '../config/ClusterReader'; export interface RouterOptions { logger: Logger; From 7112565f45bbd6ab2bd05e7e0a5d85e85b4e613f Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Thu, 21 Jan 2021 12:32:02 +0200 Subject: [PATCH 3/9] Added support for multiple consumed topics --- app-config.yaml | 14 ++-- .../kafka-backend/src/service/router.test.ts | 48 +++++++------ .../ConsumerGroupOffsets.tsx | 19 ++++-- .../useConsumerGroupsForEntity.test.tsx | 68 +++++++++++++++++-- .../useConsumerGroupsForEntity.ts | 24 +++++-- ...useConsumerGroupsOffsetsForEntity.test.tsx | 16 +++-- .../useConsumerGroupsOffsetsForEntity.ts | 25 ++++--- 7 files changed, 158 insertions(+), 56 deletions(-) diff --git a/app-config.yaml b/app-config.yaml index 02afb8ac41..7017b0eee8 100644 --- a/app-config.yaml +++ b/app-config.yaml @@ -99,6 +99,13 @@ kubernetes: - 'config' clusters: [] +kafka: + clientId: backstage + clusters: + - name: cluster + brokers: + - localhost:9092 + integrations: github: - host: github.com @@ -225,6 +232,8 @@ catalog: # Backstage example groups and users - type: file target: ../catalog-model/examples/acme-corp.yaml + - type: file + target: ../../local.yaml scaffolder: github: @@ -372,8 +381,3 @@ homepage: timezone: 'Asia/Tokyo' pagerduty: eventsBaseUrl: 'https://events.pagerduty.com/v2' - -kafka: - clientId: backstage - brokers: - - localhost:9092 diff --git a/plugins/kafka-backend/src/service/router.test.ts b/plugins/kafka-backend/src/service/router.test.ts index 400729627b..42cd92ddfb 100644 --- a/plugins/kafka-backend/src/service/router.test.ts +++ b/plugins/kafka-backend/src/service/router.test.ts @@ -51,7 +51,7 @@ describe('router', () => { jest.resetAllMocks(); }); - describe('get /:clusterId/consumer/:consumerId/offsets', () => { + describe('get /consumers/clusterId/:consumerId/offsets', () => { it('returns topic and group offsets', async () => { const topic1Offsets = [ { id: 1, offset: '500' }, @@ -82,7 +82,7 @@ describe('router', () => { .calledWith('hey') .mockResolvedValue(groupOffsets); - const response = await request(app).get('/prod/consumer/hey/offsets'); + const response = await request(app).get('/consumers/prod/hey/offsets'); expect(response.status).toEqual(200); expect(response.body.consumerId).toEqual('hey'); @@ -114,7 +114,7 @@ describe('router', () => { it('handles internal error correctly', async () => { prodKafkaApi.fetchGroupOffsets.mockRejectedValue(Error('oh no')); - const response = await request(app).get('/prod/consumer/hey/offsets'); + const response = await request(app).get('/consumers/prod/hey/offsets'); expect(response.status).toEqual(500); }); @@ -147,29 +147,35 @@ describe('router', () => { .calledWith('hey') .mockResolvedValue(groupDevOffsets); - const prodResponse = await request(app).get('/prod/consumer/hey/offsets'); - const devResponse = await request(app).get('/dev/consumer/hey/offsets'); + const prodResponse = await request(app).get( + '/consumers/prod/hey/offsets', + ); + const devResponse = await request(app).get('/consumers/dev/hey/offsets'); expect(prodResponse.status).toEqual(200); expect(prodResponse.body.consumerId).toEqual('hey'); - expect(prodResponse.body.offsets).toIncludeSameMembers([ - { - topic: 'topic1', - partitionId: 1, - groupOffset: '100', - topicOffset: '500', - }, - ]); + expect(new Set(prodResponse.body.offsets)).toStrictEqual( + new Set([ + { + topic: 'topic1', + partitionId: 1, + groupOffset: '100', + topicOffset: '500', + }, + ]), + ); expect(devResponse.status).toEqual(200); expect(devResponse.body.consumerId).toEqual('hey'); - expect(devResponse.body.offsets).toIncludeSameMembers([ - { - topic: 'topic1', - partitionId: 1, - groupOffset: '567', - topicOffset: '1234', - }, - ]); + expect(new Set(devResponse.body.offsets)).toStrictEqual( + new Set([ + { + topic: 'topic1', + partitionId: 1, + groupOffset: '567', + topicOffset: '1234', + }, + ]), + ); }); }); }); diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.tsx b/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.tsx index 448d420bc0..7928393526 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.tsx +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.tsx @@ -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={ - Consumed Topics for {consumerGroup} + Consumed Topics for {consumerGroup} ({clusterId}) } @@ -98,6 +100,15 @@ export const ConsumerGroupOffsets = ({ export const KafkaTopicsForConsumer = () => { const [tableProps, { retry }] = useConsumerGroupsOffsetsForEntity(); - - return ; + return ( + + {tableProps.consumerGroupsTopics?.map(consumerGroup => ( + + ))} + + ); }; diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx index 483986d4ec..3bee4e2e88 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.test.tsx @@ -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"`, ), ); }); diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts index 0afec42870..522a6273fe 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsForEntity.ts @@ -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; }; diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx index 16bf2f7665..b7e3a29e80 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.test.tsx @@ -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 () => { diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts index 4081b59796..8cb30040e9 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/useConsumerGroupsOffsetsForEntity.ts @@ -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, From 234e7d985188a1f740125b82857c431962d0dca2 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Thu, 21 Jan 2021 12:33:51 +0200 Subject: [PATCH 4/9] Added changeset --- .changeset/fast-flowers-tickle.md | 6 ++++++ 1 file changed, 6 insertions(+) create mode 100644 .changeset/fast-flowers-tickle.md diff --git a/.changeset/fast-flowers-tickle.md b/.changeset/fast-flowers-tickle.md new file mode 100644 index 0000000000..4a27890980 --- /dev/null +++ b/.changeset/fast-flowers-tickle.md @@ -0,0 +1,6 @@ +--- +'@backstage/plugin-kafka': patch +'@backstage/plugin-kafka-backend': patch +--- + +Added support for multiple kafka clusters and multiple consumers per component. From 4f4fd9c307155df3e4a3cc4e3ea597ebd92d5252 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Thu, 21 Jan 2021 12:37:33 +0200 Subject: [PATCH 5/9] Removed debug entry added to app-config by mistake --- app-config.yaml | 2 -- 1 file changed, 2 deletions(-) diff --git a/app-config.yaml b/app-config.yaml index 7017b0eee8..c150ef5902 100644 --- a/app-config.yaml +++ b/app-config.yaml @@ -232,8 +232,6 @@ catalog: # Backstage example groups and users - type: file target: ../catalog-model/examples/acme-corp.yaml - - type: file - target: ../../local.yaml scaffolder: github: From 4b5bd5a458379f0d257357c1af883b3a7effb843 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Thu, 21 Jan 2021 13:18:06 +0200 Subject: [PATCH 6/9] Fixed spelling error --- .changeset/fast-flowers-tickle.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changeset/fast-flowers-tickle.md b/.changeset/fast-flowers-tickle.md index 4a27890980..193b9b4459 100644 --- a/.changeset/fast-flowers-tickle.md +++ b/.changeset/fast-flowers-tickle.md @@ -3,4 +3,4 @@ '@backstage/plugin-kafka-backend': patch --- -Added support for multiple kafka clusters and multiple consumers per component. +Added support for multiple Kafka clusters and multiple consumers per component. From 9488ea7df398306b7b128a62303cb634e7460d51 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Thu, 21 Jan 2021 14:07:32 +0200 Subject: [PATCH 7/9] Fixed bug in ConsumerGroupOffsets test --- .../ConsumerGroupOffsets/ConsumerGroupOffsets.test.tsx | 1 + 1 file changed, 1 insertion(+) diff --git a/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.test.tsx b/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.test.tsx index ea0f7e1ddb..9cedc5e135 100644 --- a/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.test.tsx +++ b/plugins/kafka/src/components/ConsumerGroupOffsets/ConsumerGroupOffsets.test.tsx @@ -27,6 +27,7 @@ describe('ConsumerGroupOffsets', () => { const rendered = await renderInTestApp( {}} From 965ca841f2192e43b9a2e092dd4d4121828f569e Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Tue, 26 Jan 2021 09:22:04 +0200 Subject: [PATCH 8/9] Updated changeset with breaking changes documentation --- .changeset/fast-flowers-tickle.md | 25 +++++++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/.changeset/fast-flowers-tickle.md b/.changeset/fast-flowers-tickle.md index 193b9b4459..48051c197b 100644 --- a/.changeset/fast-flowers-tickle.md +++ b/.changeset/fast-flowers-tickle.md @@ -4,3 +4,28 @@ --- Added support for multiple Kafka clusters and multiple consumers per component. +Note that this introduces several breaking changes. + +1. Configuration in `app-config.yaml` has changed to support the ability to configure multiple clusters. This means you are required to update the configs in the following way: + +```diff +kafka: + clientId: backstage +- brokers: +- - localhost:9092 ++ clusters: ++ - name: prod ++ brokers: ++ - localhost:9092 +``` + +2. Configuration of services has changed as well to support multiple clusters: + +```diff + annotations: +- kafka.apache.org/consumer-groups: consumer ++ kafka.apache.org/consumer-groups: prod/consumer +``` + +3. Kafka Backend API has changed, so querying offsets of a consumer group is now done with the following query path: + `/consumers/${clusterId}/${consumerGroup}/offsets` From 1f1b8ba77bb0741f1ed59d85a4832b480ffce543 Mon Sep 17 00:00:00 2001 From: Nir Gazit Date: Tue, 26 Jan 2021 09:23:11 +0200 Subject: [PATCH 9/9] Updated version changes to be minor --- .changeset/fast-flowers-tickle.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/.changeset/fast-flowers-tickle.md b/.changeset/fast-flowers-tickle.md index 48051c197b..3eacb6f1a0 100644 --- a/.changeset/fast-flowers-tickle.md +++ b/.changeset/fast-flowers-tickle.md @@ -1,6 +1,6 @@ --- -'@backstage/plugin-kafka': patch -'@backstage/plugin-kafka-backend': patch +'@backstage/plugin-kafka': minor +'@backstage/plugin-kafka-backend': minor --- Added support for multiple Kafka clusters and multiple consumers per component.