chore(events): update consumerSubscribeOptions name
Signed-off-by: Jonas Beck <jonas.beck@velux.com>
This commit is contained in:
+2
-2
@@ -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();
|
||||
});
|
||||
|
||||
@@ -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 }) => {
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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'),
|
||||
},
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user