【问题标题】:How to store each Spark Streaming iteration data to one RDD?如何将每个 Spark Streaming 迭代数据存储到一个 RDD 中?
【发布时间】:2017-05-29 20:46:26
【问题描述】:

我是 Spark 的新手。 我正在编写以下脚本,它接收来自 Kafka 的流,然后将其转换为 RDD。

我的目标是将每次流迭代的数据存储在内存中到一个 RDD。就像在每个循环中向列表中添加一个元素一样。

conf = SparkConf().setAppName("Application")
sc = SparkContext(conf=conf)

def joinRDDs(rdd):
     elements = rdd.collect()
     rdds = sc.parallelize(elements)
     transformed = rdds.map(lambda x: ('key', {u'name': x[1]}))

if __name__ == '__main__':
    ssc = StreamingContext(sc, 2)
    stream = KafkaUtils.createDirectStream(ssc, [topic],{"metadata.broker.list": host})
    stream.foreachRDD(joinRDDs)

我怎样才能做到这一点?

感谢您的关注

【问题讨论】:

    标签: apache-spark pyspark apache-kafka spark-streaming rdd


    【解决方案1】:

    使用 updateStatebyKey() 并根据需要传入函数。函数接受两个参数,每批都有新数据,还有你在内存中保存的历史数据。

    def countPurchasers(newValues,lastSum): 如果 lastSum 为无: 最后总和=0 返回总和(newValues,lastSum)

    updateStatebBykey(countPurchasers)

    【讨论】:

      猜你喜欢
      • 2015-11-18
      • 2015-10-18
      • 1970-01-01
      • 2020-09-22
      • 2015-06-14
      • 2017-07-28
      • 2020-06-03
      • 1970-01-01
      • 2019-03-29
      相关资源
      最近更新 更多