如何动态合并多个从Kafka多主题订阅创建的DStream?
解决Spark Streaming合并多Kafka DStream数据丢失的问题
我来帮你搞定这个问题~你当前的代码之所以只保留了两个流的数据,是因为循环里每次都把JoinedStream重新赋值成firstStream和当前循环流的合并,之前合并的结果全被覆盖了,最后自然只剩firstStream和最后一个topic的流。另外,关于创建空DStream的方法,我也会一起告诉你。
问题根源分析
你的循环逻辑存在核心问题:
for(i <- 0 to topicNameONe.length-1 ) { JoinedStream=firstStream.union(...) }
每次循环都用firstStream和当前topic的流做union,然后覆盖掉JoinedStream——相当于每次都是重新合并firstStream和当前流,之前合并的其他流数据都被丢弃了,最后结果自然只有firstStream和最后一个topic的流。
解决方案
这里有两种可靠的实现方式,你可以根据需求选择:
方式1:用空DStream初始化(适合不确定topic列表是否为空的场景)
Spark Streaming提供了直接创建空DStream的方法,我们可以先初始化一个空流,然后逐个把每个topic的流合并进去:
import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.spark.streaming.kafka010._ // 初始化空的DStream,指定泛型为ConsumerRecord[String, String] var joinedStream: DStream[ConsumerRecord[String, String]] = streamingContext.emptyDStream[ConsumerRecord[String, String]]() // 遍历所有topic,逐个合并到joinedStream中 for (topicItem <- topicNameONe) { // 解析当前topic名称 val currentTopic = topicItem.asInstanceOf[JsString].value // 创建当前topic的DStream val currentStream = KafkaUtils.createDirectStream[String, String]( streamingContext, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](Array(currentTopic), kafkaParams) ) // 合并到已有的流中 joinedStream = joinedStream.union(currentStream) }
方式2:从第一个topic的流开始初始化(更高效,适合确定topic列表非空的场景)
如果能确定topicNameONe至少有一个topic,可以直接用第一个topic的流作为初始值,然后从第二个topic开始遍历合并:
import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.spark.streaming.kafka010._ // 用第一个topic初始化joinedStream val firstTopic = topicNameONe(0).asInstanceOf[JsString].value var joinedStream: DStream[ConsumerRecord[String, String]] = KafkaUtils.createDirectStream[String, String]( streamingContext, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](Array(firstTopic), kafkaParams) ) // 从第二个topic开始遍历合并 for (i <- 1 until topicNameONe.length) { val currentTopic = topicNameONe(i).asInstanceOf[JsString].value val currentStream = KafkaUtils.createDirectStream[String, String]( streamingContext, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](Array(currentTopic), kafkaParams) ) joinedStream = joinedStream.union(currentStream) }
额外提示
如果你的topic数量非常多,多次union可能会带来一些性能开销,不过Spark Streaming对union操作的优化做得不错,只要集群资源充足,一般不会有太大问题。如果后续遇到性能瓶颈,可以考虑调整Spark的并行度参数,比如spark.streaming.kafka.maxRatePerPartition来控制每个分区的消费速率。
内容的提问来源于stack exchange,提问作者Siddesh H K
相关产品推荐
相关产品推荐

