单Broker单Topic下Kafkajs指定消费分区及跨分区生产消费实现
实现指定Partition的消息消费与发送
一、让消费者仅消费指定Partition的消息
你只需要在subscribe方法中添加partition参数,指定要消费的分区编号(比如0、1、2,对应你的3个Partition),修改后的代码如下:
const { Kafka } = require('kafkajs'); const kafka = new Kafka({ clientId: 'my-app', brokers: ['localhost:9092'] }); const consumer = kafka.consumer({ groupId: 'test-group' }); const run = async () => { await consumer.connect(); // 指定订阅topic的第0个Partition,fromBeginning表示从开头消费 await consumer.subscribe({ topic: 'Kafka_Test_Topic', partition: 0, fromBeginning: true }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { console.log(`从分区${partition}收到消息: ${message.value.toString()}`); }, }); }; run().catch(console.error);
说明:这种方式属于手动分配分区,一旦指定了具体Partition,消费组不会再自动为这个消费者分配其他分区,确保该实例只处理你指定的分区消息。
二、实现生产者向指定Partition发消息,消费者从另一个Partition消费
1. 生产者指定Partition发送消息
创建生产者实例时,在send方法的消息配置中添加partition参数,指定要发送的目标分区:
const { Kafka } = require('kafkajs'); const kafka = new Kafka({ clientId: 'producer-app', brokers: ['localhost:9092'] }); const producer = kafka.producer(); const runProducer = async () => { await producer.connect(); // 向第1个Partition发送消息 await producer.send({ topic: 'Kafka_Test_Topic', messages: [ { value: '发送到分区1的测试消息', partition: 1 // 指定目标分区 }, ], }); await producer.disconnect(); }; runProducer().catch(console.error);
2. 消费者从另一个Partition消费
直接修改消费者的subscribe配置,指定另一个分区(比如第2个)即可:
const { Kafka } = require('kafkajs'); const kafka = new Kafka({ clientId: 'consumer-app', brokers: ['localhost:9092'] }); const consumer = kafka.consumer({ groupId: 'another-test-group' }); const runConsumer = async () => { await consumer.connect(); // 订阅第2个Partition await consumer.subscribe({ topic: 'Kafka_Test_Topic', partition: 2, fromBeginning: true }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { console.log(`从分区${partition}收到消息: ${message.value.toString()}`); }, }); }; runConsumer().catch(console.error);
这样就能实现生产者往Partition 1发消息,消费者从Partition 2消费的独立流程。
内容的提问来源于stack exchange,提问作者Gaurav
相关产品推荐
相关产品推荐

