diff --git a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.test.ts b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.test.ts index c0bf94d3c8..64c53595ad 100644 --- a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.test.ts +++ b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.test.ts @@ -39,7 +39,7 @@ describe('KafkaConsumingEventPublisher', () => { consumerConfig: { groupId: 'test-group', }, - consumerSubscribeConfig: { + consumerSubscribeTopics: { topics: ['test-topic'], }, backstageTopic: 'backstage-topic', @@ -72,7 +72,7 @@ describe('KafkaConsumingEventPublisher', () => { expect(mockConsumer.connect).toHaveBeenCalled(); expect(mockConsumer.subscribe).toHaveBeenCalledWith( - kafkaConsumerConfig.consumerSubscribeConfig, + kafkaConsumerConfig.consumerSubscribeTopics, ); expect(mockConsumer.run).toHaveBeenCalled(); }); diff --git a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.ts b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.ts index 53df325af4..cc794481cf 100644 --- a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.ts +++ b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumingEventPublisher.ts @@ -29,7 +29,7 @@ type EventMetadata = EventParams['metadata']; */ export class KafkaConsumingEventPublisher { private readonly kafkaConsumer: Consumer; - private readonly consumerSubscribeOptions: ConsumerSubscribeTopics; + private readonly consumerSubscribeTopics: ConsumerSubscribeTopics; private readonly backstageTopic: string; private readonly logger: LoggerService; @@ -54,13 +54,13 @@ export class KafkaConsumingEventPublisher { config: KafkaConsumerConfig, ) { this.kafkaConsumer = kafkaClient.consumer(config.consumerConfig); - this.consumerSubscribeOptions = config.consumerSubscribeConfig; + this.consumerSubscribeTopics = config.consumerSubscribeTopics; this.backstageTopic = config.backstageTopic; const id = `events.kafka.publisher:${this.backstageTopic}`; this.logger = logger.child({ class: KafkaConsumingEventPublisher.prototype.constructor.name, groupId: config.consumerConfig.groupId, - kafkaTopics: config.consumerSubscribeConfig.topics.toString(), + kafkaTopics: config.consumerSubscribeTopics.topics.toString(), backstageTopic: config.backstageTopic, taskId: id, }); @@ -71,7 +71,7 @@ export class KafkaConsumingEventPublisher { await this.kafkaConsumer.connect(); this.logger.info('Kafka consumer connected'); - await this.kafkaConsumer.subscribe(this.consumerSubscribeOptions); + await this.kafkaConsumer.subscribe(this.consumerSubscribeTopics); await this.kafkaConsumer.run({ eachMessage: async ({ message }) => { diff --git a/plugins/events-backend-module-kafka/src/publisher/config.test.ts b/plugins/events-backend-module-kafka/src/publisher/config.test.ts index 083c206c77..6381880437 100644 --- a/plugins/events-backend-module-kafka/src/publisher/config.test.ts +++ b/plugins/events-backend-module-kafka/src/publisher/config.test.ts @@ -71,7 +71,7 @@ describe('readConfig', () => { publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.groupId, ).toEqual('my-group'); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerSubscribeConfig.topics, + publisherConfigs.kafkaConsumerConfigs[0].consumerSubscribeTopics.topics, ).toEqual(['topic-A']); }); @@ -165,7 +165,7 @@ describe('readConfig', () => { publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.groupId, ).toEqual('my-group'); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerSubscribeConfig.topics, + publisherConfigs.kafkaConsumerConfigs[0].consumerSubscribeTopics.topics, ).toEqual(['topic-A']); expect( diff --git a/plugins/events-backend-module-kafka/src/publisher/config.ts b/plugins/events-backend-module-kafka/src/publisher/config.ts index f5e4379edd..b872a2a8d9 100644 --- a/plugins/events-backend-module-kafka/src/publisher/config.ts +++ b/plugins/events-backend-module-kafka/src/publisher/config.ts @@ -22,7 +22,7 @@ import { ConsumerConfig, ConsumerSubscribeTopics, KafkaConfig } from 'kafkajs'; export interface KafkaConsumerConfig { backstageTopic: string; consumerConfig: ConsumerConfig; - consumerSubscribeConfig: ConsumerSubscribeTopics; + consumerSubscribeTopics: ConsumerSubscribeTopics; } /** @@ -73,7 +73,7 @@ export const readConfig = (config: Config): KafkaEventSourceConfig => { maxBytes: topic.getOptionalNumber('kafka.maxBytes'), maxWaitTimeInMs: topic.getOptionalNumber('kafka.maxWaitTimeInMs'), }, - consumerSubscribeConfig: { + consumerSubscribeTopics: { topics: topic.getStringArray('kafka.topics'), }, };