Hadoop 3环境下使用GCS连接器遭遇依赖错误,修复方案可行但原理存疑
Hadoop 3环境下使用GCS连接器遭遇依赖错误,修复方案可行但原理存疑
我最近在搭建一个简单的PySpark项目,目标是把DataFrame写入Google Cloud Storage(GCS)存储桶,但一直被依赖管理相关的错误卡住。试了好几个版本的GCS连接器都没解决问题,还参考文档加了额外的配置项,结果只是换了不同的错误提示。不过后来我找到了一个可行的修复方法,但实在搞不懂它的原理,想跟大家聊聊这个情况。
我的代码实现
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType def main(): # Initialize Spark session with GCS support spark = SparkSession.builder \ .appName("Basic PySpark Example") \ .config("spark.jars.packages", "com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.26") \ .config("spark.hadoop.google.cloud.auth.service.account.enable", "true") \ .config("spark.hadoop.fs.gs.auth.service.account.json.keyfile", "<key-file>.json") \ .getOrCreate() schema = StructType([ StructField("name", StringType(), True), StructField("age", IntegerType(), True), StructField("city", StringType(), True) ]) data = [ ("John", 30, "New York"), ("Alice", 25, "London"), ("Bob", 35, "Paris") ] df = spark.createDataFrame(data, schema=schema) df.show() gcs_bucket = "name" df.write.parquet(f"gs://{gcs_bucket}/data/people.parquet") df.write.csv(f"gs://{gcs_bucket}/data/people.csv") print(f"Data written to GCS bucket: {gcs_bucket}") spark.stop() if __name__ == "__main__": main()
requirements.txt配置
pyspark==3.5.5 google-cloud-storage==3.1.0
遇到的错误信息
py4j.protocol.Py4JJavaError: An error occurred while calling o49.parquet. : org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "gs" at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:3443) at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3466) at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174) at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3574) at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3521) at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:540) at org.apache.hadoop.fs.Path.getFileSystem(Path.java:365) at org.apache.spark.sql.execution.datasources.DataSource.planForWritingFileFormat(DataSource.scala:454) at org.apache.spark.sql.execution.datasources.DataSource.planForWriting(DataSource.scala:530) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:391) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:364) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:243) at org.apache.spark.sql.DataFrameWriter.parquet(DataFrameWriter.scala:802) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:75) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:52) at java.base/java.lang.reflect.Method.invoke(Method.java:580) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374) 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.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) at py4j.ClientServerConnection.run(ClientServerConnection.java:106) at java.base/java.lang.Thread.run(Thread.java:1583)
可行的修复方案
把SparkSession配置里的这一行:
.config("spark.jars.packages", "com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.26")
替换成:
.config("spark.jars", "https://storage.googleapis.com/hadoop-lib/gcs/gcs-connector-hadoop3-latest.jar")
我的疑问
为什么通过spark.jars.packages指定Maven坐标的方式会失败,而直接指定JAR包的URL就能正常工作呢?有没有大佬能帮忙解释一下背后的原因?
备注:内容来源于stack exchange,提问作者Jesus Diaz Rivero
相关产品推荐
相关产品推荐

