【问题标题】:Scala to PysparkScala 到 Pyspark
【发布时间】: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


    【解决方案1】:

    您的代码需要进行一些更改:

    • 您不能广播RDD:而是在底层“数据”上广播:
    • 然后您可以使用value() 方法在闭包内获取广播变量

    以下是您更新后的代码的近似值:

     #Create static data
        data = [('log_name','enrichment_success')])
        #Broadcast it to all nodes
        ip_classification_broadcast = sc.broadcast(data)
        #Join stream with static dataset on field log_name      
        joinedStream = kafkaStream.transform(lambda rdd:  \
            rdd.join(ip_classification_broadcast.value().get[1]))
    

    【讨论】:

    • 错误日志片段:rdd.join(ip_classification_broadcast.value().get()[log_name])) TypeError: 'list' object is not callable
    • 看起来value 是一个属性而不是一个方法——所以我更新了上面的代码:尝试.get[log_name] 而不是get()[log_name]。你需要稍微调整一下这些细节,因为我没有你的代码和测试设置。
    • 为了调试,我已经运行了join函数的内容。 ip_classification_broadcast.value().get[log_name] 导致:TypeError: 'list' object is not callable 这:ip_classification_broadcast.value[0] 导致:('log_name', 'enrichment_success') 但是当我运行 spark-submit 时,这是错误:AttributeError: 'tuple' object has no attribute 'mapValues'
    • 啊 - 这意味着我们已经接近了!您可能需要重组data,使其成为dict 而不是tuple。但在此之前 - 通过使用当前的 tuple: 以下应该编译并给你 enrichment_success 作为输出: ip_classification_broadcast.value().get[1] 。注意访问tuple 的第二个元素的[1]。稍后您将希望能够通过密钥log_name 访问:这将需要重组。
    • CASE 1: ip_classification_broadcast = sc.broadcast(data) joinedStream = kafkaStream.transform(lambda rdd: rdd.join(ip_classification_broadcast.value[0])) ERROR AttributeError: 'tuple' object has no attribute 'mapValues'
    猜你喜欢
    • 1970-01-01
    • 2021-01-31
    • 2018-03-19
    • 1970-01-01
    • 2018-02-03
    • 2017-01-20
    • 2019-06-09
    • 2018-04-03
    • 2020-05-08
    相关资源
    最近更新 更多