【发布时间】: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