【发布时间】:2018-04-19 10:58:49
【问题描述】:
我正在使用 Apache Spark 开发流媒体应用程序。该应用程序通过订阅名为 sensor 的 Kafka 主题来接收传感器数据。该应用程序的目的是过滤传感器数据,对其进行转换并将其发布回名为 people 的不同 Kafka 主题以供其他消费者使用。主题people 中的消息必须与它们到达主题sensor 的顺序相同。因此,我目前在 Kafka 中只使用一个分区。
这是我的代码:
val myStream = KafkaUtils.createDirectStream[K, V](streamingContext, PreferConsistent, Subscribe[K, V](topics, consumerConfig))
def process(record: (RDD[ConsumerRecord[String, String]], Time)): Unit = record match {
case (rdd, time) if !rdd.isEmpty =>
// More Code...
// Filter RDD, transform to JSON, build Seq[People]...
// In the end, I have: Dataset[People]
// Publish to Kafka topic 'people'
case _ =>
}
myStream.foreachRDD((x, y) => process((x, y)))
今天,我问了一个问题,关于如何在 Spark 中实现正确的排序,将其转换为我的People 数据结构。
answer 表示将 Spark 与单个分区一起使用是不明智的,这实际上可能是一个设计缺陷:
除非你有一个单独的分区(然后你不会使用 Spark,对吗?)订单...
我现在想知道是否可以改进应用程序的整体设计(更改 map-reduce 流程),或者 Spark 是否不适合我的用例。
【问题讨论】:
标签: scala apache-spark apache-kafka spark-streaming