| 
 | 1 | +import {  | 
 | 2 | +  Consumer,  | 
 | 3 | +  EachMessagePayload,  | 
 | 4 | +  Kafka,  | 
 | 5 | +  Partitioners,  | 
 | 6 | +  ProducerRecord  | 
 | 7 | +} from 'kafkajs';  | 
 | 8 | +import config from '../utils/config';  | 
 | 9 | +import { Xml } from './ediSacha/sachaToXML';  | 
 | 10 | + | 
 | 11 | +const kafka = new Kafka({  | 
 | 12 | +  clientId: 'client-id',  | 
 | 13 | +  brokers: [config.kafka.url]  | 
 | 14 | +});  | 
 | 15 | + | 
 | 16 | +export const createConsumer = async (): Promise<Consumer> => {  | 
 | 17 | +  const consumer = kafka.consumer({  | 
 | 18 | +    groupId: 'consumer-group'  | 
 | 19 | +  });  | 
 | 20 | + | 
 | 21 | +  try {  | 
 | 22 | +    await consumer.connect();  | 
 | 23 | +    await consumer.subscribe({  | 
 | 24 | +      topic: config.kafka.topicRAI  | 
 | 25 | +    });  | 
 | 26 | + | 
 | 27 | +    await consumer.run({  | 
 | 28 | +      eachMessage: async (messagePayload: EachMessagePayload) => {  | 
 | 29 | +        const { topic, partition, message } = messagePayload;  | 
 | 30 | +        const prefix = `${topic}[${partition} | ${message.offset}] / ${message.timestamp}`;  | 
 | 31 | +        console.log(`- ${prefix} ${message.key}#${message.value}`);  | 
 | 32 | +      }  | 
 | 33 | +    });  | 
 | 34 | +  } catch (error) {  | 
 | 35 | +    console.log('Error: ', error);  | 
 | 36 | +  }  | 
 | 37 | + | 
 | 38 | +  return consumer;  | 
 | 39 | +};  | 
 | 40 | + | 
 | 41 | +export const sendMessage = async (message: Xml): Promise<void> => {  | 
 | 42 | +  const producer = kafka.producer({  | 
 | 43 | +    createPartitioner: Partitioners.DefaultPartitioner  | 
 | 44 | +  });  | 
 | 45 | + | 
 | 46 | +  const topicMessages: ProducerRecord = {  | 
 | 47 | +    topic: config.kafka.topicDAI,  | 
 | 48 | +    messages: [{ value: message }]  | 
 | 49 | +  };  | 
 | 50 | +  await producer.connect();  | 
 | 51 | +  await producer.send(topicMessages);  | 
 | 52 | +  await producer.disconnect();  | 
 | 53 | +};  | 
0 commit comments