【发布时间】:2019-01-24 10:30:12
【问题描述】:
我正在尝试在 Dstream 和静态 RDD 之间执行连接。
PySpark
#Create static data
ip_classification_rdd = sc.parallelize([('log_name','enrichment_success')])
#Broadcast it to all nodes
ip_classification_rdd_broadcast = sc.broadcast(ip_classification_rdd)
#Join stream with static dataset on field log_name
joinedStream = kafkaStream.transform(lambda rdd: rdd.join(ip_classification_rdd[log_name]))
我得到了这个例外: "您似乎正在尝试广播一个 RDD 或从一个 "
引用一个 RDD斯卡拉
但是,这里有人有同样的要求:How to join a DStream with a non-stream file?
这就是解决方案:
val vdpJoinedGeo = goodIPsFltrBI.flatMap{ip => geoDataBC.value.get(ip).map(data=> (ip,data)}
在 Pyspark 中这个等价物是什么?
【问题讨论】:
标签: scala apache-spark pyspark spark-streaming dstream