【问题标题】:Passing function to spark that reads S3 file using pyspark将函数传递给使用 pyspark 读取 S3 文件的 spark
【发布时间】:2020-04-21 05:53:17
【问题描述】:

我在 s3 中有 GB 的数据,并尝试通过引用以下 Link 来读取我的代码时带来并行性。

我使用下面的代码作为示例,但是当我运行时,它运行到以下错误:

对此的任何帮助都深表感谢,因为我对火花很陌生。

编辑:我必须使用并行性读取我的 s3 文件,这在任何帖子中都没有解释。标记重复的人请先阅读问题。

PicklingError:无法序列化对象:异常:您似乎正试图从广播变量、操作或转换中引用 SparkContext。 SparkContext 只能在驱动程序上使用,不能在它在工作人员上运行的代码中使用。有关详细信息,请参阅 SPARK-5063。

class telco_cn:
    def __init__(self, sc):
        self.sc = sc

    def decode_module(msg):
        df=spark.read.json(msg)
        return df

    def consumer_input(self, sc, k_topic):
        a = sc.parallelize(['s3://bucket1/1575158401-51e09537-0ce5-c775-6beb-fd1b0a568e15.json'])
        d = a.map(lambda x: telco_cn.decode_module(x)).collect()
        print (d)
if __name__ == "__main__":
    cn = telco_cn(sc)
    cn.consumer_input(sc, '')

【问题讨论】:

标签: python apache-spark pyspark bigdata


【解决方案1】:

您正试图从 RDD 上的 map 操作中调用 spark.read.json。由于此映射操作将在 Spark 的 executor/worker 节点上执行,因此您无法在映射中引用 SparkContext/SparkSession 变量(在 Spark 驱动程序上定义)。这就是错误消息试图告诉您的内容。

为什么不直接打电话给df=spark.read.json('s3://bucket1/1575158401-51e09537-0ce5-c775-6beb-fd1b0a568e15.json')

【讨论】:

  • 那么如果我不使用地图,我将如何实现并行性? @查理
  • spark.read.json是并行操作,并行度由输入文件的数量定义(一般来说每个输入文件1个任务)。
猜你喜欢
  • 2020-02-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-06-09
  • 1970-01-01
  • 1970-01-01
  • 2018-08-19
  • 2015-09-11
相关资源
最近更新 更多