【发布时间】: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