【发布时间】:2018-03-15 20:05:47
【问题描述】:
我正在使用 Spark Streaming,并且我开发了以下 Spark Streaming 应用程序:
从 Kafka 接收器 (RDD1) 创建一个 DStream,从 HTTP 请求 (RDD2) 创建另一个。
我的问题是,我只想使用 RDD1 中的第一个元素并在我的 RDD2 中使用它,并且此代码在 spark 流式传输 (.first()) 中不起作用如何使用 spark 流式传输 1.6 获得相同的结果
代码:
firstLineRDD = kvs.map(lambda x : x[0], x[1].split('\n')[0], x[2])
dateRDD = firstLineRDD.map(lambda x : (datetime.datetime.fromtimestamp(float(x[0])/1000000),x[1],x[2]))
dayAggRDD = dateRDD.map(lambda x : (x[0],x[1],x[2]))
daily_date, sys , metric = dayAggRDD.first()
dataTSRDD = sc.parallelize(apiRequest(sys,metric,getDailyDate(daily_date)))
【问题讨论】:
-
是
RDD还是DStreams ?? -
@massg 我已经使用 DStreams 修改了我的帖子
-
azelix,我还是没看到你指的
DStreams。
标签: apache-spark pyspark spark-streaming