【问题标题】:How to filter dstream using transform operation and external RDD?如何使用转换操作和外部 RDD 过滤 dstream?
【发布时间】:2015-09-02 03:18:33
【问题描述】:

我在Transformations on DStreams转换操作部分中描述的类似用例中使用了transform方法:

spamInfoRDD = sc.pickleFile(...) # RDD containing spam information
# join data stream with spam information to do data cleaning
cleanedDStream = wordCounts.transform(lambda rdd: rdd.join(spamInfoRDD).filter(...))

我的代码如下:

sc = SparkContext("local[4]", "myapp")
ssc = StreamingContext(sc, 5)
ssc.checkpoint('hdfs://localhost:9000/user/spark/checkpoint/')
lines = ssc.socketTextStream("localhost", 9999)
counts = lines.flatMap(lambda line: line.split(" "))\
              .map(lambda word: (word, 1))\
              .reduceByKey(lambda a, b: a+b)
filter_rdd = sc.parallelize([(u'A', 1), (u'B', 1)], 2)
filtered_count = counts.transform(
    lambda rdd: rdd.join(filter_rdd).filter(lambda k, (v1, v2): v1 and not v2)
)
filtered_count.pprint()
ssc.start()
ssc.awaitTermination()

但我收到以下错误

您似乎正在尝试广播 RDD 或从操作或转换中引用 RDD。 RDD 转换和操作只能由驱动程序调用,不能在其他转换内部调用;例如,rdd1.map(lambda x: rdd2.values.count() * x) 无效,因为值转换和计数操作无法在 rdd1.map 转换内部执行。有关详细信息,请参阅 SPARK-5063。

我应该如何使用我的外部 RDD 从 dstream 中过滤元素?

【问题讨论】:

  • 你有答案吗

标签: apache-spark spark-streaming pyspark


【解决方案1】:

Spark 文档示例和您的代码之间的区别在于 ssc.checkpoint() 的使用。

虽然您提供的特定代码示例在没有检查点的情况下也可以工作,但我想您实际上需要它。但是将外部 RDD 引入检查点 DStream 范围的概念可能无效:当从检查点恢复时,外部 RDD 可能已更改。

我尝试检查点外部 RDD,但我也没有运气。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-06-29
    • 1970-01-01
    • 1970-01-01
    • 2021-01-23
    • 2019-02-07
    • 2017-04-05
    • 1970-01-01
    相关资源
    最近更新 更多