【问题标题】:Pyspark: how to read a .csv file in google bucket?Pyspark:如何读取谷歌存储桶中的 .csv 文件?
【发布时间】:2021-04-01 17:20:10
【问题描述】:

我有一些文件存储在谷歌存储桶中。这些是我建议的here 的设置。

spark = SparkSession.builder.\
        master("local[*]").\
        appName("TestApp").\
        config("spark.serializer", KryoSerializer.getName).\
        config("spark.jars", "/usr/local/.sdkman/candidates/spark/2.4.4/jars/gcs-connector-hadoop2-2.1.1.jar").\
        config("spark.kryo.registrator", GeoSparkKryoRegistrator.getName).\
        getOrCreate()
#Recommended settings for using GeoSpark
spark.conf.set("spark.driver.memory", 6)
spark.conf.set("spark.network.timeout", 1000)
#spark.conf.set("spark.driver.maxResultSize", 5)
spark.conf.set

spark._jsc.hadoopConfiguration().set('fs.gs.impl', 'com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem')
# This is required if you are using service account and set true, 
spark._jsc.hadoopConfiguration().set('fs.gs.auth.service.account.enable', 'false')
spark._jsc.hadoopConfiguration().set('google.cloud.auth.service.account.json.keyfile', "myJson.json")



path = 'mBucket-c892b51f8579.json'
os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = path
client = storage.Client()
name = 'https://console.cloud.google.com/storage/browser/myBucket/'
bucket_id = 'myBucket'
bucket = client.get_bucket(bucket_id)

我可以使用以下内容简单地阅读它们:

df = pd.read_csv('gs://myBucket/myFile.csv.gz', compression='gzip')
df.head()

    time_zone_name           province_short
0   America/Chicago              US.TX
1   America/Chicago              US.TX
2   America/Los_Angeles          US.CA
3   America/Chicago              US.TX
4   America/Los_Angeles          US.CA

我正在尝试使用 pyspark 读取相同的文件

myTable = spark.read.format("csv").schema(schema).load('gs://myBucket/myFile.csv.gz', compression='gzip')

但我收到以下错误

Py4JJavaError: An error occurred while calling o257.load.
: java.lang.NoClassDefFoundError: Could not initialize class com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem
    at java.lang.Class.forName0(Native Method)
    at java.lang.Class.forName(Class.java:348)
    at org.apache.hadoop.conf.Configuration.getClassByNameOrNull(Configuration.java:2134)
    at org.apache.hadoop.conf.Configuration.getClassByName(Configuration.java:2099)
    at org.apache.hadoop.conf.Configuration.getClass(Configuration.java:2193)
    at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:2654)
    at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2667)
    at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:94)
    at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:2703)
    at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:2685)
    at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:373)
    at org.apache.hadoop.fs.Path.getFileSystem(Path.java:295)
    at org.apache.spark.sql.execution.streaming.FileStreamSink$.hasMetadata(FileStreamSink.scala:45)
    at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:332)
    at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:223)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:211)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:178)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
    at py4j.Gateway.invoke(Gateway.java:282)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.GatewayConnection.run(GatewayConnection.java:238)
    at java.lang.Thread.run(Thread.java:748)

【问题讨论】:

  • @user10938362 这有点复杂。你有什么建议或一些代码要分享吗?
  • 对于测试,如果您在本地有文件,这是否有效?看起来它可能与从 URL 获取它有关,可能会出现身份验证问题。在本地驱动器上尝试 myFile.csv.gz
  • @user10938362 你说的答案并没有解决问题
  • 您似乎缺少连接器。尝试按照here 中的步骤操作。确保下载的jar 文件包含com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem 类文件。

标签: python apache-spark google-cloud-platform pyspark


【解决方案1】:

欢迎来到 hadoop 依赖地狱!

1.使用包而不是 jars

您的配置基本上是正确的,但是当您将 gcs-connector 添加为本地 jar 时,您还需要手动确保其所有依赖项在 JVM 类路径中可用。

将连接器添加为一个包并让 spark 处理依赖项通常更容易,因此使用 config('spark.jars.packages', 'com.google.cloud.bigdataoss:gcs-connector:hadoop2-2.1.1') 而不是 config('spark.jars.packages', 'com.google.cloud.bigdataoss:gcs-connector:hadoop2-2.1.1')

2.管理 ivy2 依赖项解析问题

当您按照上述方式进行操作时,spark 可能会抱怨由于 maven(用于发布)和 ivy2(由 spark 用于依赖项解析)之间的分辨率差异,它无法下载某些依赖项。

您通常可以通过简单地使用spark.jars.excludes 让 spark 忽略未解决的依赖项来解决此问题,因此添加一个新的配置行,例如 config('spark.jars.excludes','androidx.annotation:annotation,org.slf4j:slf4j-api')

3.管理类路径冲突

完成此操作后,SparkSession 将启动,但文件系统仍将失败,因为 pyspark 的标准分发打包了旧版本的 guava 库,该库未实现 gcs-connector 所依赖的 API。

您需要通过使用以下配置config('spark.driver.userClassPathFirst','true')config('spark.executor.userClassPathFirst','true') 确保 gcs-connector 将首先找到其预期版本

4.管理依赖冲突

现在您可能认为一切正常,但实际上并非如此,因为默认的 pyspark 发行版包含 2.7.3 版的 hadoop 库,但 gcs-connector 2.1.1 版仅依赖于 2.8+ 的 API。

现在你的选择是:

  • 将自定义构建的 spark 与较新的 hadoop(或没有内置 hadoop 库的包)一起使用
  • 使用旧版本的 gcs-connector(版本 1.9.17 可以正常工作)

5.终于有一个工作配置了

假设您想坚持使用 PyPi 或 Anaconda 最新的 pyspark 发行版,以下配置应该可以按预期工作。

我只包含了 gcs 相关配置,将 Hadoop 配置直接移动到 spark 配置中,并假设您正确设置了 GOOGLE_APPLICATION_CREDENTIALS:

from pyspark.sql import SparkSession

spark = SparkSession.builder.\
        master("local[*]").\
        appName("TestApp").\
        config('spark.jars.packages', 
               'com.google.cloud.bigdataoss:gcs-connector:hadoop2-1.9.17').\
        config('spark.jars.excludes',
               'javax.jms:jms,com.sun.jdmk:jmxtools,com.sun.jmx:jmxri').\
        config('spark.driver.userClassPathFirst','true').\
        config('spark.executor.userClassPathFirst','true').\
        config('spark.hadoop.fs.gs.impl',
               'com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem').\
        config('spark.hadoop.fs.gs.auth.service.account.enable', 'false').\
        getOrCreate()

请注意,gcs-connector 版本 1.9.17 有一组与 2.1.1 不同的排除项,因为为什么不...

PS:您还需要确保您使用的是 Java 1.8 JVM,因为 Spark 2.4 不适用于较新的 JVM。

【讨论】:

    【解决方案2】:

    除了@rluta 的出色回答,您还可以通过将番石榴罐子专门放在 extraClassPath 中来替换 userClassPathFirst 行:

    spark.driver.extraClassPath=/root/.ivy2/jars/com.google.guava_guava-27.0.1-jre.jar:/root/.ivy2/jars/com.google.guava_failureaccess-1.0.1.jar:/root/.ivy2/jars/com.google.guava_listenablefuture-9999.0-empty-to-avoid-conflict-with-guava.jar
    spark.executor.extraClassPath=/root/.ivy2/jars/com.google.guava_guava-27.0.1-jre.jar:/root/.ivy2/jars/com.google.guava_failureaccess-1.0.1.jar:/root/.ivy2/jars/com.google.guava_listenablefuture-9999.0-empty-to-avoid-conflict-with-guava.jar
    

    这有点麻烦,因为您需要确切的本地 ivy2 路径,但您也可以将 jar 下载/复制到更永久的位置。

    但是,这减少了其他潜在的依赖冲突,例如 livy,如果 gcs-connector 的 jackson 依赖位于类路径前面,则会抛出 java.lang.NoClassDefFoundError: org.apache.livy.shaded.json4s.jackson.Json4sModule

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-10-25
      • 2016-12-29
      • 2021-12-23
      • 2020-12-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-11-11
      相关资源
      最近更新 更多