如何使用Spark Java读取GCS中的CSV文件?
解决Spark读取GCS文件时报"No FileSystem for scheme: gs"错误
问题重现
使用Spark Java项目读取GCS上的CSV文件,代码如下:
SparkConf conf = new SparkConf().setMaster("local[*]").setAppName("Hello world"); SparkSession sparkSession = SparkSession.builder().config(conf).getOrCreate(); Dataset<Row> dataset = sparkSession.read() .option("header", true) .option("sep", ",") .option("delimiter", "\"") .csv("gs://abc/WDC_age.csv");
运行时抛出错误:
Exception in thread "main" java.io.IOException: No FileSystem for scheme: gs
解决方案
1. 添加GCS Hadoop客户端依赖
Spark需要依赖Google Cloud Storage的Hadoop客户端库来识别gs协议,根据你的Spark版本添加对应依赖:
Maven(适配Spark 3.x + Hadoop 3.x)
<dependency> <groupId>com.google.cloud.bigdataoss</groupId> <artifactId>gcs-connector-hadoop3</artifactId> <version>2.2.0</version> </dependency>
Gradle(适配Spark 3.x + Hadoop 3.x)
implementation 'com.google.cloud.bigdataoss:gcs-connector-hadoop3:2.2.0'
注意:若Spark基于Hadoop 2.x开发,替换为
gcs-connector-hadoop2,选择对应兼容版本即可。
2. 配置Spark的FileSystem实现类
在SparkConf中指定gs协议对应的FileSystem实现类:
SparkConf conf = new SparkConf() .setMaster("local[*]") .setAppName("Hello world") .set("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem");
3. 修正CSV读取参数错误
代码中option("delimiter", "\"")属于配置错误:delimiter用于指定字段分隔符(已通过sep设为逗号),双引号作为字段引用符应使用quote参数配置,或直接省略(默认值即为双引号)。修正后的读取逻辑:
Dataset<Row> dataset = sparkSession.read() .option("header", true) .option("sep", ",") .option("quote", "\"") // 可选,默认值即为双引号 .csv("gs://abc/WDC_age.csv");
4. 本地运行的认证配置(可选)
若在本地环境运行,需配置GCS服务账号认证:
- 下载服务账号密钥JSON文件
- 通过系统属性指定密钥路径:
System.setProperty("GOOGLE_APPLICATION_CREDENTIALS", "/path/to/your/service-account-key.json");
或在SparkConf中配置:
conf.set("fs.gs.auth.service.account.json.keyfile", "/path/to/your/service-account-key.json");
完整修正代码
import org.apache.spark.SparkConf; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; public class GcsCsvReader { public static void main(String[] args) { // 本地运行需配置认证 System.setProperty("GOOGLE_APPLICATION_CREDENTIALS", "/path/to/service-account-key.json"); SparkConf conf = new SparkConf() .setMaster("local[*]") .setAppName("GCS CSV Reader") .set("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem"); SparkSession sparkSession = SparkSession.builder().config(conf).getOrCreate(); Dataset<Row> dataset = sparkSession.read() .option("header", true) .option("sep", ",") .csv("gs://abc/WDC_age.csv"); // 验证读取结果 dataset.show(); sparkSession.stop(); } }
内容的提问来源于stack exchange,提问作者Faisal Khan
相关产品推荐
相关产品推荐

