Spark Structured Streaming关联rate与csv流批次耗时过长如何优化
问题根因
你当前遇到的长耗时核心是流-流关联缺少必要的时间约束配置,Spark Structured Streaming为了保证无界流关联的结果正确性,默认会永久保留两个流的所有历史状态做匹配,没有时间约束的情况下会一直等待可能的匹配数据,同时你还存在资源配置、触发策略的错误,共同导致了长耗时。
可落地的优化方案
1. 修正流-流关联的约束配置(核心优化)
流-流内连接必须满足两个条件才能启用状态自动清理,避免无限等待:
- 两个流都需要预先定义事件时间字段并设置watermark
- 关联条件中必须显式指定两个流事件时间的匹配范围,让Spark可以判断哪些旧状态可以安全删除
你当前的csv流没有自带事件时间,可以用数据摄入时间代替,修改后的代码示例:
// csv流新增摄入时间和watermark val cvsStream=readFromCSVFile(tmpPath.toString) .withColumn("ingestTime", current_timestamp()) .withWatermark("ingestTime", "5 seconds") // 关联时增加时间范围约束,不要在join之后再加watermark val cvsStream2 = cvsStream.as("csv").join(rate.as("counter"), expr(""" csv.id = counter.value AND counter.timestamp >= csv.ingestTime AND counter.timestamp <= csv.ingestTime + interval 10 seconds """) )
2. 调整运行资源配置
你当前使用的local[1]模式只有1个工作线程,而流处理至少需要2个线程(1个用于数据摄入,1个用于任务计算),单核心会直接导致摄入和计算任务抢资源阻塞,把master配置修改为local[*]或者local[2]即可。
3. 调整触发和分区配置
- 你当前设置的
Trigger.ProcessingTime(10)是10毫秒触发一次,过高的触发频率会产生大量空批次调度开销,测试场景下改为1秒触发即可:Trigger.ProcessingTime(1000) - rate流的
numPartitions设置为和核心数匹配即可,不要设置远大于核心数的分区,会增加不必要的调度开销。
4. 可选优化:改用流-静态关联
如果你的csv数据是不会频繁更新的维度数据,直接用静态读csv代替流读csv,Spark会自动使用广播join优化,性能会提升几个量级:
// 静态读取csv数据 val csvStatic = spark.read .format("csv") .option("header", true) .schema(schema) .option("delimiter", ",") .load(tmpPath.toString) // 直接和rate流关联,不需要watermark和时间约束 val result = rate.join(csvStatic, expr("value = id"))
内容的提问来源于stack exchange,提问作者Eljah
相关产品推荐
相关产品推荐

