【问题标题】:Kafka (Re-)Joining group stuck with more than 2 topics卡夫卡(重新)加入小组,主题超过 2 个
【发布时间】:2017-09-12 20:17:34
【问题描述】:

我正在开发一个使用 Kafka 作为消息发布/订阅工具的系统。

数据由 scala 脚本生成:

val kafkaParams = new Properties()
    kafkaParams.put("bootstrap.servers", "localhost:9092")
    kafkaParams.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    kafkaParams.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    kafkaParams.put("group.id", "test_luca")

    //kafka producer
    val producer = new KafkaProducer[String, String](kafkaParams)

    //Source list
    val s1 = new java.util.Timer()
    val tasks1 = new java.util.TimerTask {
        def run() = {
            val date = new java.util.Date
            val date2 = date.getTime()
            val send = ""+ date2 + ", 45.1234, 12.5432, 4.5, 3.0"
            val data = new ProducerRecord[String,String]("topic_s1", send)
            producer.send(data)
        }
    }
    s1.schedule(tasks1, 1000L, 1000L)

    val s2 = new java.util.Timer()
    val tasks2 = new java.util.TimerTask {
        def run() = {
            val date = new java.util.Date
            val date2 = date.getTime()
            val send = ""+ date2 + ", 1.111, 9.999, 10.4, 10.0"
            val data = new ProducerRecord[String,String]("topic_s2", send)
            producer.send(data)
        }
    }
    s2.schedule(tasks2, 2000L, 2000L)

我需要在某些特定情况下测试 kafka 的性能。在一种情况下,我有另一个脚本使用来自主题“topic_s1”和“topic_s2”的数据,详细说明它们,然后生成具有不同主题(topic_s1b 和 topic_s2b)的新数据。随后,这些详细的数据由 Apache Spark Streaming 脚本使用。

如果我省略消费者/生产者脚本(我只有 1 个 Kafka 生产者,有 2 个主题和 Spark 脚本),一切正常。

如果我使用完整配置(1 个带有 2 个主题的 kafka 生产者,使用来自 kafka 生产者的数据的“中间件”脚本,详细说明它们并使用新主题生成新数据,1 个使用新主题使用数据的 spark 脚本) Spark Streaming 脚本卡在INFO AbstractCoordinator: (Re-)joining group test_luca

我在本地运行所有东西,我不会修改 kafka 和 zookeeper 配置。

有什么建议吗?

更新:火花脚本:

val sparkConf = new SparkConf().setAppName("SparkScript").set("spark.driver.allowMultipleContexts", "true").setMaster("local[2]")
val sc = new SparkContext(sparkConf)

val ssc = new StreamingContext(sc, Seconds(4))

case class Thema(name: String, metadata: JObject)
case class Tempo(unit: String, count: Int, metadata: JObject)
case class Spatio(unit: String, metadata: JObject)
case class Stt(spatial: Spatio, temporal: Tempo, thematic: Thema)
case class Location(latitude: Double, longitude: Double, name: String)

case class Data(location: Location, timestamp: Long, measurement: Int, unit: String, accuracy: Double)
case class Sensor(sensor_name: String, start_date: String, end_date: String, data_schema: Array[String], data: Data, stt: Stt)


case class Datas(location: Location, timestamp: Long, measurement: Int, unit: String, accuracy: Double)
case class Sensor2(sensor_name: String, start_date: String, end_date: String, data_schema: Array[String], data: Datas, stt: Stt)


val kafkaParams = Map[String, Object](
    "bootstrap.servers" -> "localhost:9092",
    "key.deserializer" -> classOf[StringDeserializer].getCanonicalName,
    "value.deserializer" -> classOf[StringDeserializer].getCanonicalName,
    "group.id" -> "test_luca",
    "auto.offset.reset" -> "latest",
    "enable.auto.commit" -> (false: java.lang.Boolean)
)

val topics1 = Array("topics1")
val topics2 = Array("topics2")

val stream = KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics1, kafkaParams))
val stream2 = KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics2, kafkaParams))

val s1 = stream.map(record => {
  implicit val formats = DefaultFormats
  parse(record.value).extract[Sensor]
}
)
val s2 = stream2.map(record => {
  implicit val formats = DefaultFormats
  parse(record.value).extract[Sensor2]
}
)

val f1 = s1.map { x => x.sensor_name }
f1.print()
val f2 = s2.map { x => x.sensor_name }
f2.print()

谢谢 卢卡

【问题讨论】:

  • 请显示您的 spark 流式处理脚本代码。
  • @GuangshengZuo 我已经上传了 Spark 脚本

标签: scala apache-kafka spark-streaming


【解决方案1】:

也许您应该更改 spark 流脚本的 group.id。我猜您的“中间件”脚本的使用者与您的火花流脚本的使用者具有相同的 group.id。然后可怕的事情就会发生。

在kafka中,消费者组是主题的真正订阅者,组中的消费者只是一个分裂工作者,所以在你的情况下,你应该在中间件脚本消费者和火花流脚本消费者中使用不同的group.id。

在您第一次尝试没有中间脚本的情况下,它只是因为这个而起作用。

【讨论】:

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