如何使用直接流在Kafka Spark Streaming中指定使用者组

Fai*_*qui 4 java apache-kafka apache-spark spark-streaming kafka-consumer-api

如何使用直接流API为kafka spark流指定使用者组ID.

HashMap<String, String> kafkaParams = new HashMap<String, String>();
kafkaParams.put("metadata.broker.list", brokers);
kafkaParams.put("auto.offset.reset", "largest");
kafkaParams.put("group.id", "app1");

    JavaPairInputDStream<String, String> messages = KafkaUtils.createDirectStream(
            jssc, 
            String.class, 
            String.class,
            StringDecoder.class, 
            StringDecoder.class, 
            kafkaParams, 
            topicsSet
    );
Run Code Online (Sandbox Code Playgroud)

虽然我已经指定配置不确定是否遗漏了什么.使用spark1.3

kafkaParams.put("group.id", "app1");
Run Code Online (Sandbox Code Playgroud)

C4s*_*tor 6

直接流API使用低级Kafka API,因此无论如何都不使用消费者组.如果要将消费者组与Spark Streaming一起使用,则必须使用基于接收器的API.

doc中提供了完整的详细信息!


sec*_*ree 5

createDirectStreaminspark-streaming-kafka-0-8不支持组模式,因为它使用的是低级 Kafka API。

spark-streaming-kafka-0-10支持组模式。

消费者配置

在 0.9.0.0 中,我们引入了新的 Java 消费者作为旧的基于 Scala 的简单和高级消费者的替代品。新旧消费者的配置如下所述。

在 中New Consumer Configs,它有group.id项目。

Spark Streaming integration for Kafka 0.10是使用新的API。https://spark.apache.org/docs/2.1.1/streaming-kafka-0-10-integration.html

Kafka 0.10 的 Spark Streaming 集成在设计上类似于 0.8 Direct Stream 方法。它提供了简单的并行性、Kafka 分区和 Spark 分区之间的 1:1 对应关系,以及对偏移量和元数据的访问。但是,由于较新的集成使用了新的 Kafka 消费者 API 而不是简单的 API,因此在用法上存在显着差异。

我已经在 中测试了组模式spark-streaming-kafka-0-10,它确实有效。