Spark Streaming按groupId写入Redshift无法并行执行,如何解决?
问题根因
你遇到的串行问题核心有两点:
- Scala并行集合
par的默认并行度由Driver进程的可用CPU核心数决定,YARN Client模式下Driver通常默认分配1核,所以实际并行度为1,看起来就是串行运行。 - 默认FIFO调度器下,所有Spark作业按提交顺序排队执行,即使你有多个线程提交作业,也会被调度器串行执行。
解决方案
方案一:调整调度模式+自定义并行集合并行度(改动最小)
适合要求每个groupId写入为独立事务、单独提交的业务场景。
- 调整Spark调度模式为公平调度,支持多作业并行运行,提交作业时新增配置:
--conf spark.scheduler.mode=FAIR
代码中配置示例:
val spark = SparkSession.builder() .config("spark.scheduler.mode", "FAIR") // 其余原有配置 .getOrCreate()
- 手动设置并行集合的并行度,覆盖默认的Driver核数限制,并行度建议不超过可用Executor总核数(你的环境3个Executor每个2核,建议设3~5即可),修改后代码如下:
import scala.concurrent.forkjoin.ForkJoinPool import scala.collection.parallel.ForkJoinTaskSupport inputDstream.foreachRDD { eventRdd: RDD[Event] => ... // Convert eventRdd to eventDF val groupIds = eventDF.select("group_id").distinct.collect.flatMap(_.toSeq) val parGroupIds = groupIds.par // 自定义并行度为3,可根据实际资源调整 parGroupIds.tasksupport = new ForkJoinTaskSupport(new ForkJoinPool(3)) parGroupIds.foreach{ groupId => val teventDF = eventDF.where($"group_id" <=> groupId) val teventDFWithVersion = teventDF.withColumn("schema_id", lit(version)) teventDFWithVersion.write .format("io.github.spark_redshift_community.spark.redshift") .options(opts) .mode("Append") .save() } }
方案二:使用Spark原生分布式写入(更推荐,性能更高)
无特殊事务要求优先选择该方案,完全利用Spark分布式计算能力,避免Driver端并行度瓶颈,代码如下:
inputDstream.foreachRDD { eventRdd: RDD[Event] => ... // Convert eventRdd to eventDF val eventDFWithVersion = eventDF.withColumn("schema_id", lit(version)) // 按group_id重分区,确保同一个group_id的数据落到同一个分区 val repartitionedDF = eventDFWithVersion.repartition($"group_id") // 直接写入Redshift,连接器会分布式并行处理所有分区数据 repartitionedDF.write .format("io.github.spark_redshift_community.spark.redshift") .options(opts) .mode("Append") .save() }
该方案不需要修改调度模式,也不需要使用并行集合,性能和稳定性远高于Driver端多线程提交的方案。
内容的提问来源于stack exchange,提问作者jn5047
相关产品推荐
相关产品推荐

