import kafka, { Consumer } from 'kafka-node' import logger from '../logger/Logger' import BaseProvider from '../core/BaseProvider' export default class KafkaProvider extends BaseProvider { connectDataSource(): void { const kafkaHost = this.params.host const client = new kafka.KafkaClient({ kafkaHost }) const consumer = new Consumer(client, this.params.topics, { autoCommit: true, groupId: "Stream-Computation" }) let offset = new kafka.Offset(client); consumer.on("message", (message) => { try { this.onMessage(message.topic, JSON.parse(message.value.toString())) } catch (err) { logger.info(err) } }) consumer.on("error", (err) => { logger.info(err) }) consumer.on('offsetOutOfRange', function (topic) { logger.info("------------- offsetOutOfRange ------------"); topic.maxNum = 2; offset.fetch([topic], function (err, offsets) { logger.info(offsets); var min = Math.min.apply(null, offsets[topic.topic][topic.partition]); consumer.setOffset(topic.topic, topic.partition, min); }); }); logger.info("kafka连接成功") } }