【问题标题】:Kafka Node High Level Producer writes to even partitions onlyKafka Node High Level Producer 仅写入偶数分区
【发布时间】:2019-02-17 09:46:48
【问题描述】:

我正在使用 Kafka Node 库,并测试高级生产者。

我创建了一个包含 10 个分区的主题“HLPTestInput”,并编写了一个函数以每秒生成一次。

生产者写入分区 0、2、4、6 和 8,但不写入奇数分区。

奇怪的是,当我从这个主题消费并生产到第二个主题“HLPTestInputFromConsumer”时,它有 5 个分区,消息被写入所有分区。

是否有我遗漏的配置?

const kafka = require('kafka-node'),
    HighLevelProducer = kafka.HighLevelProducer,
    ConsumerGroup = kafka.ConsumerGroup,
    client = new kafka.KafkaClient({kafkaHost: 'smc-dev.silverbolt.lab:9092'}),
    producer = new HighLevelProducer(client),
    consumer = new ConsumerGroup(
        {
          kafkaHost: 'smc-dev.silverbolt.lab:9092',
            groupId: 'testGroup'
        },
        'HLPTestInput'
    );

let index = 0;
setInterval(() => {
    producer.send([{
        topic: 'HLPTestInput',
        messages: [index]
    }], (err, data) => {
        console.log('produced', data);
    });
    index++;
}, 1000);

consumer.on('message', (message) => {
    console.log('consumed', message);
    producer.send([{
        topic: 'HLPTestInputFromConsumer',
        messages: [message]
    }], (err, data) => {
        console.log('produced to secondary', data);
    });
});

【问题讨论】:

    标签: javascript apache-kafka node-kafka


    【解决方案1】:

    我不太确定,但可能是因为你使用同一个制作人来写两个不同的主题。由于 HighLevelProducer 使用循环来编写。因此,假设您的生产者在“HLPTestInput”主题中写入,然后您将时间间隔设置为 1000,因此在此期间,您的消费者收到消息,现在您的生产者在“HLPTestInputFromConsumer”主题中写入。

    所以您的生产者在其分区 0、2、4 中写入“HLPTestInput”主题...

    和“HLPTestInputFromConsumer”主题在其第 1,3,5 部分 ...

    所以我建议尝试创建另一个生产者。那么它应该可以正常工作。

    试试下面的代码:

    const kafka = require('kafka-node'),
        HighLevelProducer = kafka.HighLevelProducer,
        ConsumerGroup = kafka.ConsumerGroup,
        client = new kafka.KafkaClient({kafkaHost: 'smc-dev.silverbolt.lab:9092'}),
        client1 = new kafka.KafkaClient({kafkaHost: 'smc-dev.silverbolt.lab:9092'}),
        producer = new HighLevelProducer(client),
        producer1 = new HighLevelProducer(client1),
        consumer = new ConsumerGroup(
           {
              kafkaHost: 'smc-dev.silverbolt.lab:9092',
               groupId: 'testGroup'
            },
            'HLPTestInput'
        );
    let index = 0;
        setInterval(() => {
        producer.send([{
            topic: 'HLPTestInput',
            messages: [index]
        }], (err, data) => {
            console.log('produced', data);
        });
       index++;
    }, 1000);
    
    consumer.on('message', (message) => {
        console.log('consumed', message);
        producer1.send([{
            topic: 'HLPTestInputFromConsumer',
            messages: [message]
        }], (err, data) => {
            console.log('produced to secondary', data);
        });
    });
    

    【讨论】:

      猜你喜欢
      • 2014-04-12
      • 2019-04-23
      • 1970-01-01
      • 2022-01-15
      • 2016-01-21
      • 1970-01-01
      • 1970-01-01
      • 2016-06-23
      • 1970-01-01
      相关资源
      最近更新 更多