【问题标题】:Design flaw in Spark Streaming App by using only one partition?Spark Streaming App 仅使用一个分区的设计缺陷?
【发布时间】: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


    【解决方案1】:

    虽然这主要是基于意见的您使用的工具是为:

    • 容错,
    • 分布式,
    • 平行,
    • 处理,没有具体的订单保证

    解决问题定义为:

    • 顺序,
    • 非分布式,
    • 有严格的订单保证,
    • 可能会破坏容错(由于大量数据放置在单个执行器上)。

    地点:

    • 来自容错队列的单线程消费者

    完全够用了

    所以主观上来说这里有一个严重的设计缺陷。

    【讨论】:

      【解决方案2】:

      在您的情况下,Kafka 不是正确的选择。 Kafka 仅维护分区内消息的总顺序。 Kafka 的并行性或可扩展性完全取决于特定主题的分区数。缺陷完全在于设计。

      如果你真的想保持秩序,你可以有一个纪元 数据中的时间戳,一旦转换数据,您就可以排序 数据并存储它。

      【讨论】:

        猜你喜欢
        • 2018-09-26
        • 1970-01-01
        • 2016-01-04
        • 1970-01-01
        • 1970-01-01
        • 2014-08-19
        • 1970-01-01
        相关资源
        最近更新 更多