【发布时间】: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