基于函数式编程范式,用kafka-node实现Node.js的Pub/Sub封装模块
基于kafka-node的函数式Kafka Pub/Sub封装方案
针对你的需求,下面是完全基于函数式编程范式封装的KafkaLib模块,通过高阶函数、纯函数设计实现简洁的publish和subscribe调用接口,同时保持逻辑的可组合性和无副作用。
核心实现代码(KafkaLib.js)
const kafka = require('kafka-node'); // 创建Kafka客户端连接的纯函数 const createKafkaClient = (kafkaHost) => new kafka.KafkaClient({ kafkaHost }); // 异步创建生产者实例,返回Promise const createProducer = (client) => { const producer = new kafka.Producer(client); return new Promise((resolve, reject) => { producer.on('ready', () => resolve(producer)); producer.on('error', reject); }); }; // 创建消费者实例的纯函数 const createConsumer = (client, topic, groupId, config = {}) => { return new kafka.Consumer( client, [{ topic }], { groupId, autoCommit: true, ...config } ); }; // 柯里化的发布函数:先绑定Kafka地址,再返回具体的发布操作 const publish = (kafkaHost) => async (topic, messages) => { const client = createKafkaClient(kafkaHost); const producer = await createProducer(client); return new Promise((resolve, reject) => { producer.send([{ topic, messages }], (err, data) => { client.close(); // 发送完成后主动关闭连接 err ? reject(err) : resolve(data); }); }); }; // 柯里化的订阅函数:先绑定Kafka地址和消费组,再返回订阅操作,同时返回取消订阅函数 const subscribe = (kafkaHost, groupId) => (topic, messageHandler, consumerConfig = {}) => { const client = createKafkaClient(kafkaHost); const consumer = createConsumer(client, topic, groupId, consumerConfig); // 绑定消息处理逻辑 consumer.on('message', messageHandler); consumer.on('error', (err) => console.error('Kafka consumer error:', err)); // 返回取消订阅的清理函数,符合函数式资源管理思路 return () => { consumer.close(true, () => client.close()); }; }; module.exports = { publish, subscribe };
使用示例
生产者端调用
const { publish } = require('./KafkaLib'); // 预先绑定Kafka服务地址,生成可复用的发布函数 const publishUserEvent = publish('localhost:9092'); // 发布单条/多条消息 (async () => { try { // 发布单条字符串消息 const result = await publishUserEvent('user-events', 'user: alice registered'); console.log('消息发布成功:', result); // 发布多条消息 await publishUserEvent('user-events', ['user: bob logged in', 'user: bob updated profile']); } catch (err) { console.error('发布失败:', err); } })();
消费者端调用
const { subscribe } = require('./KafkaLib'); // 预先绑定Kafka地址和消费组ID,生成可复用的订阅函数 const subscribeToUserEvents = subscribe('localhost:9092', 'user-events-consumer-group'); // 订阅主题并处理消息,同时获取取消订阅函数 const unsubscribe = subscribeToUserEvents('user-events', (message) => { console.log('收到消息:', message.value); }); // 如需停止订阅,调用unsubscribe()即可清理资源 // setTimeout(() => { // console.log('停止订阅'); // unsubscribe(); // }, 60000);
设计说明
- 函数式范式体现:采用柯里化设计,将配置与业务操作分离;所有创建实例的函数都是纯函数,输入相同则输出相同;通过返回清理函数管理资源,避免副作用泄漏。
- 简洁性:上层调用只需关注业务逻辑(发布/订阅主题、处理消息),无需关心Kafka客户端的创建、连接、销毁细节。
- 可扩展性:
createConsumer支持传入自定义配置,可覆盖默认的autoCommit等参数;如需批量发布、自定义分区等功能,可扩展publish函数的参数。
内容的提问来源于stack exchange,提问作者Subhasish Biswasray
相关产品推荐
相关产品推荐

