【问题标题】:MessageHandler in KafkaUtils010 SparkStreamingKafkaUtils 010 Spark Streaming 中的 MessageHandler
【发布时间】:2017-06-09 19:26:43
【问题描述】:

我想按主题分组或在申请时知道消息来自哪个主题:

val stream = KafkaUtils.createDirectStream[String, String](
    ssc,
    PreferConsistent, 
    Subscribe[String, String](
      Array(topicConfig.srcTopic),
      kafkaParameters(BOOTSTRAP_SERVERS,"kafka_test_group_id))
    )
  )

但是,在最新的 API 中,kafka010 似乎不像以前的版本那样支持消息处理程序。关于如何获得主题的任何想法?

我的目标是从 N 个主题中消费处理它们(根据主题以不同的方式),然后以 1:1 的主题映射将其推送回另一个 N 个主题:

SrcTopicA--> Process --> DstTopicA
SrcTopicB--> Process --> DstTopicB
SrcTopicC--> Process --> DstTopicC

但是有一些属性需要共享(变化很大,所以不可能使用广播变量)。因此,所有主题都需要在同一个 Spark 作业中使用。

【问题讨论】:

    标签: scala apache-spark apache-kafka spark-streaming


    【解决方案1】:

    当你在 0.10 中使用 createDirectStream 时,你会得到一个 ConsumerRecord。该记录具有topic 值。您可以创建一个主题和值的元组:

    val stream: InputDStream[ConsumerRecord[String, String]] = 
      KafkaUtils.createDirectStream[String, String](
        streamingContext,
        PreferConsistent,
        Subscribe[String, String](topics, kafkaParams)
      )
    
    val res: DStream[(String, String)] = stream.map(record => (record.topic(), record.value()))
    

    【讨论】:

      【解决方案2】:

      您可以像这样使用topic 过滤流:

      stream.filter(cr => cr.topic)
      

      【讨论】:

        猜你喜欢
        • 2017-02-21
        • 1970-01-01
        • 2016-12-19
        • 2015-02-26
        • 2020-03-19
        • 1970-01-01
        • 1970-01-01
        • 2020-09-12
        • 1970-01-01
        相关资源
        最近更新 更多