【发布时间】: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, '')
【问题讨论】:
-
上面的例子我已经看过了。在这里,他们只是试图将一些值附加到字典中,而我的用例非常不同。我正在尝试读取 df 中的 json。 @blackbishop
标签: python apache-spark pyspark bigdata