【问题标题】:SQL Stored Procedure to Scala/Spark StreamingSQL 存储过程到 Scala/Spark 流
【发布时间】:2017-01-06 03:19:46
【问题描述】:

我目前正在将一个主要用 SQL 存储过程编写的古老系统转移到 Scala 以在 Spark 上运行。存储过程是每天/每周/每月/每年在“请求”对象上运行一次的批处理作业,可能需要数小时才能运行。

出于几个原因,我们正在将系统更改为流模型(Spark Streaming)。

在旧系统中,很多逻辑都是通过join语句来执行的,其中大量的Requests与许多表进行join。

一种解决方案是本质上采用相同的 SQl 代码并将其移植到 Spark SQL 语句中,然后这些语句将在请求的“微批量”上运行。然而,这意味着我们仍在执行大量的连接语句,据我所知,这些语句在 Spark SQL 中效率低下。

我的第二个想法是采用业务逻辑并编写代码,就好像我们只需要处理一个请求(即,如果您有 10 个应用程序,而不是处理所有应用程序使用连接,您可以像处理单个请求一样进行编程)。然后,我将获取微批次的请求并通过逻辑处理(即 Requests.map(r => RequestLogic.execute(r)))映射它们。

类似于以下示例代码:

case class Request(id: Int, typeId: Int, value: Long)

def CreateStreamingContext(sparkConf: SparkConf, streamDuration: Duration,
                             storageLevel: StorageLevel = StorageLevel.MEMORY_ONLY): StreamingContext = {

    sparkConf.set(SparkArgumentKeys.MaxCores, (partitionCount * 2).toString)
    val ssc = new StreamingContext(sparkConf, streamDuration)
    ssc.checkpoint(checkpointDir)

    val stream = EventHubsUtils.createUnionStream(ssc, hubParams, storageLevel)
    stream.checkpoint(streamDuration)

    stream.map(x => Request(x(1), x(2), x(3)))
      .map(r => RequestLogic.execute(r))

    ssc
}

我想弄清楚:

1) 哪个扩展性更好。
2) 各有什么优缺点。

我是 Scala/Spark 的新手,并试图找出最好的方法。我不确定这是否足够的信息,如果需要,我会尝试提供更多详细信息。

【问题讨论】:

    标签: scala apache-spark spark-streaming spark-dataframe


    【解决方案1】:

    有趣的问题,答案取决于数据的形状。我想假设两种情况:

    • 首先,您有很多 Request 数据,并且您希望将它们加入到数量相对较少的主数据中。

    • 其次,Request 数据的数据量和要加入的数据一样大,超过了集群的 RAM。

    在第一种情况下,您可以建议 Spark(原则上它也应该能够自动确定)使用名为 BroadcastHashJoin 的东西。策略是将小表广播给每个 Spark 工作人员,并将其与较大 RDD 中的每个元素连接起来。将有超过 2 个非空分区,因此 Spark 会在两个以上工作节点时运行得更快。在合同中,ShuffleHashJoin 将把所有的行都放进去,并用密钥打乱它们。整个表只有 2 个非空分区,向作业添加更多工作节点也无济于事。因此,能够执行BroadcastHashJoin 可以确保以最少的工作量实现可扩展性。有关更多详细信息,请参阅 DataBricks 中的链接笔记本。

    在第二种情况下,您的策略不是一个坏主意。我们在预处理数据以加入外部 KV 存储(如 RocksDBAWS DynamoDB)方面有很好的经验,然后进行查找并以流方式加入。但是,尽管即使在小型集群上也能够将此过程扩展到真正的海量数据集,但其性能和工作量比纯 Spark 内存方法要高得多。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-07-12
      • 2018-12-06
      • 2018-05-08
      相关资源
      最近更新 更多