【问题标题】:Spark-streaming for task parallelization用于任务并行化的 Spark-streaming
【发布时间】:2016-10-02 07:05:11
【问题描述】:

我正在设计一个具有以下流程的系统:

  1. 通过网络下载提要文件(基于行)
  2. 将元素解析为对象
  3. 过滤无效/不必要的对象
  4. 在部分元素上执行阻塞 IO(HTTP 请求)
  5. 保存到数据库

我一直在考虑使用 Spark-streaming 实现系统,主要用于任务并行化、资源管理、容错等。

但我不确定这是火花流的正确用例,因为我不只是将它用于度量和数据处理。 另外我不确定 Spark-streaming 如何处理阻塞的 IO 任务。

Spark-streaming 是否适合这种用例?或者我应该寻找其他技术/框架?

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    Spark 的核心是一个通用并行计算框架。 Spark Streaming 添加了一个抽象来支持使用微批处理的流处理。 我们当然可以在 Spark Streaming 上实现这样的用例。

    要“扇出” I/O 操作,我们需要在两个级别上确保正确的并行度:

    • 首先,将数据均匀分布在分区之间: 数据的初始分区将取决于使用的流式传输源。对于这个用例,看起来custom receiver 可能是要走的路。收到批处理后,我们可能需要使用dstream.repartition(n) 来处理更大数量的分区,这些分区应该大致匹配为作业分配的执行器数量的 2-3 倍。

    • Spark 为每个执行的任务使用 1 个内核(可配置)。每个分区执行任务。这假设我们的任务是 CPU 密集型的并且需要完整的 CPU。为了优化阻塞 I/O 的执行,我们希望为许多操作多路复用该内核。我们通过直接对分区进行操作并使用经典的并发编程来并行化我们的工作来做到这一点。

    鉴于feedLinesDstream 的原始流,我们可以这样: (* 在 Scala 中。Java 版本应该类似,但 LOC 是 x 倍)

    val feedLinesDstream = ??? // the original dstream of feed lines
    val parsedElements = feedLinesDstream.map(parseLine)
    val validElements = parsedElements.filter(isValid _)
    val distributedElements = validElements.repartition(n) // n = 2 to 3 x #of executors
    
    // multiplex execution at the level of each partition
    val data =  distributedElements.mapPartitions{ iter =>
       implicit executionContext = ??? // obtain a thread pool for execution
       val futures = iter.map(elem => Future(ioOperation(elem)))
       // traverse the future resulting in a future collection of results
       val res = Future.sequence(future) 
       Await.result(res, timeout)
    }
    data.saveToCassandra(keyspace, table)
    

    【讨论】:

      【解决方案2】:

      Spark-streaming 是否适合这种用例?或者我应该看看 其他技术/框架?

      在考虑使用 Spark 时,您应该问自己几个问题:

      1. 我的应用程序在当前状态下的规模是多少,将来会发展到什么程度? (Spark 通常适用于每秒会发生数百万个进程的大数据应用程序

      2. 我首选哪种语言? (Spark 可以用 Java、Scala、PythonR 实现)

      3. 我将使用什么数据库? (Apache Spark 等技术通常使用 HBase 等大型数据库结构实现)

      我也不确定 Spark-streaming 如何处理阻塞的 IO 任务。

      Stack Overflow 上已经有一个answer,关于在 Scala 中使用 Spark 阻塞 IO 任务。它应该给你一个开始,但要回答这个问题,是的是有可能的。

      最后,阅读文档很重要,你可以找到 Spark 的权利here

      【讨论】:

      • 系统应该支持每秒数万,使用Java和Cassandra。如果 Spark 不是完全正确的选择,我仍然想享受任务并行化、高可用性和可扩展性,还有什么更合适的吗?
      • 嗯,任务量相当大!继续寻找火花。特别是因为您使用的是 Java 和 Cassandra(Java 是 Spark 的一等公民)。请让我更新。
      猜你喜欢
      • 2021-04-01
      • 2018-01-13
      • 2019-09-29
      • 2020-05-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-12-18
      • 2017-06-19
      相关资源
      最近更新 更多