You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

单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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 04:50:02