【问题标题】:Multiple Consumer setup for Single Producer with 4 partitions Kafka Java具有 4 个分区 Kafka Java 的单个生产者的多个消费者设置
【发布时间】:2017-10-20 16:51:16
【问题描述】:

我创建了一个具有 4 个分区的简单生产者,现在想在一个消费者组中创建 4 个消费者来消费来自每个分区的数据。 我该怎么做?

消费者代码

    public class KafkaConsumer {
    static List<String> list = new ArrayList<String>();
    public  static DataFrame reqFieldBOIDDf;
    public static DataFrame df ;
    static SparkConf conf = new SparkConf()

                     .setAppName("kafka-sandbox")

                     .setMaster("local[*]")

                     .set("spark.cassandra.connection.host","localhost"); //for cassandra

             static JavaSparkContext sc = new JavaSparkContext(conf);

    private static long lastOffset;
            public static void main(String[] str) throws InterruptedException {
               execute();
            }
            private static void execute() throws InterruptedException {

                    KafkaConsumer<String, String> consumer = createConsumer();       consumer.subscribe(Arrays.asList("KafkaConsumerTopic6"));
                    processRecords(consumer);
                    System.out.println("Inside execute");
            }
            private static KafkaConsumer<String, String> createConsumer() {
                    Properties props = new Properties();
                    props.put("bootstrap.servers", "localhost:9092");
                    String consumeGroup = "cg1";
                    props.put("group.id", consumeGroup);
                    props.put("enable.auto.commit", "true");
                    props.put("auto.commit.interval.ms", "101");
                    props.put("max.partition.fetch.bytes", "1035");
                    props.put("heartbeat.interval.ms", "3000");
                    props.put("session.timeout.ms", "6001");props.put("max.poll.records","500");
props.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer");
                   props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
                    return new KafkaConsumer<String, String>(props);
            }
            private static void processRecords(KafkaConsumer<String, String> consumer) throws InterruptedException  {
                  while (true){
                           ConsumerRecords<String, String> records = consumer.poll(1000);
                            lastOffset = 0;
                            for (ConsumerRecord<String, String> record : records) {
                            lastOffset = record.offset();
                            list.add(record.value());
                             }
                   }
            }

一小段代码会很有帮助。

提前致谢。

【问题讨论】:

    标签: java apache-kafka kafka-consumer-api


    【解决方案1】:

    只需像启动 4 个生产者一样启动 4 个消费者。

    确保它们都具有相同的group.id 设置,并让它们都订阅主题(或主题,如果它是 4 个主题和 1 个分区)。由于它们将在同一个组中,Kafka 会自动为每个分区分配一个消费者。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-09-11
      • 2020-06-17
      • 2019-02-04
      • 1970-01-01
      • 2018-03-07
      • 1970-01-01
      相关资源
      最近更新 更多