【问题标题】:Batch lookup data for Spark streamingSpark 流的批量查找数据
【发布时间】:2016-05-31 16:51:48
【问题描述】:

我需要从 HDFS 上的文件中查找 Spark 流作业中的一些数据 此数据由批处理作业每天获取一次。
这样的任务是否有“设计模式”?

  • 如何在 a
    之后立即重新加载内存中的数据(哈希图) 每日更新?
  • 如何在查找数据时持续提供流式作业
    被取走?

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    一种可能的方法是删除本地数据结构并改用有状态流。假设您有一个名为mainStream 的主数据流:

    val mainStream: DStream[T] = ???
    

    接下来您可以创建另一个读取查找数据的流:

    val lookupStream: DStream[(K, V)] = ???
    

    还有一个可以用来更新状态的简单函数

    def update(
      current: Seq[V],  // A sequence of values for a given key in the current batch
      prev: Option[V]   // Value for a given key from in the previous state
    ): Option[V] = { 
      current
        .headOption    // If current batch is not empty take first element 
        .orElse(prev)  // If it is empty (None) take previous state
     }
    

    这两个部分可以用来创建状态:

    val state = lookup.updateStateByKey(update)
    

    剩下的就是键入mainStream并连接数据:

    def toPair(t: T): (K, T) = ???
    
    mainStream.map(toPair).leftOuterJoin(state)
    

    虽然从性能的角度来看,这可能不是最佳的,但它利用了现有的架构,让您无需手动处理失效或故障恢复。

    【讨论】:

    • 谢谢!我需要每天重新获取所有查找数据。所以我正在考虑清除现有的并应用新的。 updateStateByKey 在这种情况下会起作用吗?不清楚的是我如何每天一次将数据读取到 DStream 中。至于加入你的意思是我把查找流中的所有查找数据与主流中的记录连接起来?
    • 流媒体只关心传入的数据。如果每天将新数据推送到查找流一次,它将每天更新一次。加入我的意思正是我在示例中使用的代码类型。据您了解,这是您想要的那种操作。
    • 我正在使用您的建议作为参考,它工作正常。但是,如果您能对更新功能进行一些解释,那就太好了。我正在尝试将其重写为mapWithState
    • 当然,现在更容易理解了吗?也可以查看stackoverflow.com/q/35563876/1560062
    • 感谢您的更新。我想对此有进一步的问题,创建了一个新问题:stackoverflow.com/questions/37550054/…
    猜你喜欢
    • 2021-02-06
    • 1970-01-01
    • 2021-12-03
    • 2015-07-14
    • 2023-03-21
    • 1970-01-01
    • 1970-01-01
    • 2019-01-14
    • 2017-02-10
    相关资源
    最近更新 更多