【问题标题】:Spark streaming 24X7 with updateStateByKey Issue带有 updateStateByKey 问题的 Spark 24X7 流式传输
【发布时间】:2018-08-10 16:05:22
【问题描述】:

我正在运行 24/7 的火花流并使用 updateStateByKey 是否可以 24/7 运行火花流? 如果是,updateStateByKey 不会变大,如何处理? 当我们 24/7 运行时,我们是否必须定期重置/删除 updateStateByKey,如果不是如何以及何时重置它? 还是 Spark 以分布式方式处理?如何动态创建内存/存储。

当 updateStateByKey 增长时出现以下错误

Array out of bound exception

Exception while deleting local spark dir: /var/folders/3j/9hjkw0890sx_qg9yvzlvg64cf5626b/T/spark-local-20141026101251-cfb4
java.io.IOException: Failed to delete: /var/folders/3j/9hjkw0890sx_qg9yvzlvg64cf5626b/T/spark-local-20141026101251-cfb4

如何处理。如果有任何文档,请指出我?我完全被卡住了,非常感谢您的帮助.. 感谢您的宝贵时间

【问题讨论】:

    标签: spark-streaming


    【解决方案1】:

    在 Java 中使用 Optional.absent() 和在 Scala 中使用 None 来删除键。可以在http://blog.cloudera.com/blog/2014/11/how-to-do-near-real-time-sessionization-with-spark-streaming-and-apache-hadoop/ 找到工作示例。

    【讨论】:

    • 非常感谢我会尝试
    【解决方案2】:

    使用 None 更新密钥会将其从 spark 中删除。如果想保留key一定的时间,可以给key附加一个过期时间,每批次检查一次。

    例如,这里是按分钟计算记录的代码:

    val counts = lines.map(line => (currentMinute, 1))
    val countsWithState = counts updateStateByKey { (values: Seq[Int], state: Option[Int]) =>
      if (values.isEmpty) { // every key will be iterated, even if there's no record in this batch
        println("values is empty")
        None // this will remove the key from spark
      } else {
        println("values size " + values.size)
        Some(state.sum + values.sum)
      }
    }
    

    【讨论】:

    • 能否请您分享如何设置无。例如,如果我每 1 小时应用一次 groupby,现在是 2 小时,到 3 小时 2 小时数据不会改变如何设置无 2 小时密钥并将其保存到外部文件以供将来使用?任何帮助都非常感谢..
    • Function2, Optional, Optional> updateFunction =new Function2, Optional, Optional>() { @Override public Optional call(List values, Optional state) { Double newSum = state.or(0D); if(values.isEmpty()){ System.out.println("空值");返回 ScalaLang.none(); }else{ for (double i : values) { newSum += i;返回 Optional.of(newSum); }}};不工作,请咨询
    • 我返回选项 当我尝试设置 scalaLang.none() 它显示错误
    【解决方案3】:

    pyspark : updateStateByKey(self, updateFunc, numPartitions=None, initialRDD=None)

    返回一个新的“状态”DStream,其中每个键的状态通过应用更新 键的先前状态和键的新值的给定函数。

    @param updateFunc:状态更新函数。如果此函数返回 None,则 相应的状态键值对将被淘汰。

    updateFunc方法返回None,状态键值对删除;

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-03-06
      • 2018-05-27
      • 1970-01-01
      • 2020-07-17
      • 1970-01-01
      • 2016-08-25
      • 2016-09-21
      • 1970-01-01
      相关资源
      最近更新 更多