【问题标题】:Could not find group info of kafka when using spark streaming使用火花流时找不到 kafka 的组信息
【发布时间】:2018-08-23 06:55:23
【问题描述】:

我有以下简单的火花流进度,它使用来自 kafka 主题 test 的消息,组 ID 为 feature1,并且只打印结果。但是,当我运行bin/kafka-consumer-groups.sh --bootstrap-server zookeeper-1:9092 --list 列出所有组时,没有feature1 或任何包含feature1 的内容。有什么问题?

我的spark版本是2.1.2,kafka版本是2.12-2.0.0,zookeeper版本是3.4.13。我在这里https://github.com/yahoo/kafka-manager/issues/207 发现了一些与之相关的问题,但我不知道我的问题与问题有关。

# coding=utf8

import sys
import datetime
import time

from pyspark import SparkContext, SparkConf
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils

if __name__ == "__main__":
    spark_conf = SparkConf()
    spark_conf.set('spark.streaming.kafka.maxRatePerPartition', 1)
    sc = SparkContext("local[2]", "NetworkWordCount", conf=spark_conf)
    ssc = StreamingContext(sc, 10)

    # Create a DStream that will connect to hostname:port, like localhost:9999
    kafka_params = {
        "bootstrap.servers":"zookeeper-1:9092",
        "group.id":"feature1",
        "auto.offset.reset":"smallest",
        "session.timeout.ms":"60000",
        "request.timeout.ms":"100000",
    }
    lines = KafkaUtils.createDirectStream(ssc, ["test"], kafka_params)
    # lines = KafkaUtils.createStream(ssc, 'zookeeper-1:2181', 'feature1', {'new-one':1})
    lines.pprint()

    ssc.start()
    ssc.awaitTermination()

组列表的输出如下,sudo 没有任何改变。

console-consumer-9215
console-consumer-41888
console-consumer-32417
console-consumer-35073
console-consumer-66656

还有一个奇怪的现象,feature1出现在zookeeper的/consumers目录中,而console-consumer-*groups没有出现。

【问题讨论】:

  • 刚刚试过这个命令,它对我有用,你能发布你得到的回应吗?您是否看到此 kafka 服务器上连接了其他组 ID? (也许添加 sudo)
  • 我已经发布了我收到的回复。 console-consumer-* 组 id 由 bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning 创建。正常工作。 @AbhiskekN

标签: python apache-kafka spark-streaming


【解决方案1】:

下面的代码 sn-p 是 kafka 脚本在后端运行以获取消费者组的内容。试试这个来打印消费者组,看看你的组是否正在打印(注意“过滤器”),它在 Scala 中。

卡夫卡版本:1.0.0 , 斯卡拉版本:2.12.0

import kafka.admin.AdminClient

def main(args: Array[String]): Unit = {

    val props = new Properties()
    props.put("bootstrap.servers","<kafka-bootstrap>:9092")

    AdminClient.create(props).listAllConsumerGroupsFlattened().map(_.groupId).filter(_.contains("mx-")).mkString(";").split(";").foreach(println(_))

  }

【讨论】:

  • 我该如何运行它?我不熟悉scala。我已经安装了 sbt。能详细描述一下吗?
猜你喜欢
  • 2017-04-27
  • 1970-01-01
  • 2018-12-22
  • 2020-03-10
  • 2017-03-27
  • 2017-01-01
  • 2019-03-02
  • 2018-07-27
  • 2019-10-11
相关资源
最近更新 更多