本地Hadoop Spark工作流写入Avro至GCS实操配置咨询
GCS路径识别配置
不需要做块设备挂载或者HDFS联邦这类映射操作,只要让Hadoop生态能识别gs://路径协议即可,配置步骤如下:
- 在集群所有节点的Hadoop公共依赖目录(一般是
$HADOOP_HOME/share/hadoop/common/lib/)放入和集群Hadoop版本匹配的GCS connector jar包,比如Hadoop 3.x版本对应gcs-connector-hadoop3-*版本的包,不要跨大版本混用,否则会报类找不到或者方法不存在的错误。 - 修改所有节点的
core-site.xml,添加GCS文件系统实现和认证配置:- 配置
fs.gs.impl值为com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem - 配置
fs.AbstractFileSystem.gs.impl值为com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS - 提前申请好有GCS桶读写权限的服务账号,将JSON密钥文件同步到所有节点的固定只读路径,配置
google.cloud.auth.service.account.json.keyfile指向该密钥文件路径
- 配置
- 配置完成后在任意节点执行
hadoop fs -ls gs://<你的目标桶名>/,能正常列出桶内文件就说明配置生效,后续Spark、Oozie都可以直接用gs://前缀的路径读写GCS,不需要额外映射。
别图省事用GCS FUSE把桶挂载到本地节点当磁盘写,性能比原生Hadoop FileSystem实现差30%以上,大任务量下容易出现文件句柄泄漏、写入超时的问题,直接走官方connector是生产环境验证过的最稳方案。
Parquet任务改Avro写入GCS的调整方法
代码层面改动非常少,核心是替换输出格式和路径:
- 首先确认Spark环境有对应版本的spark-avro依赖,Spark 2.4之后Avro作为官方独立组件维护,spark-avro的版本必须和你集群运行的Spark大版本完全一致。
- 找到原来写Parquet到HDFS的代码,做两处修改即可:
- 将写入格式从parquet替换为avro:把原有的
.parquet(outputPath)写法改为.format("avro").save(outputPath),如果引入了Spark Avro的隐式转换,也可以直接用.avro(outputPath) - 将原有的
hdfs://开头的输出路径,替换为gs://<桶名>/<业务存储路径>格式的GCS路径
- 将写入格式从parquet替换为avro:把原有的
- 兼容性提示:建议写入时显式指定Avro schema,不要依赖Spark自动推断,避免生成的Avro schema和下游消费端的解析规则不匹配,同时可以指定压缩格式为snappy,和原Parquet的压缩策略保持一致,控制存储成本。
- 代码示例:
// 原有Parquet写HDFS逻辑 // val outputPath = "hdfs:///dwd/order_daily/dt=20240520" // df.write.mode("overwrite").parquet(outputPath) // 修改后Avro写GCS逻辑 val outputPath = "gs://dwd-core-bucket/order_daily/dt=20240520" df.write .mode("overwrite") .format("avro") .option("compression", "snappy") .option("avroSchema", avroSchemaStr) // 显式传入提前定义好的Avro schema .save(outputPath)
spark-submit提交注意事项
- 依赖校验:如果集群没有全局部署spark-avro和GCS connector包,提交时必须通过
--jars参数传入对应版本的jar包,或者直接打入作业胖包,注意排查依赖冲突,尤其是不要混入不同Scala版本、不同Spark版本的avro相关包,否则会报序列化异常。 - 临时路径修改:如果要完全弃用HDFS承载数据,必须把Hadoop全局临时路径从HDFS改到GCS,否则Spark Shuffle之外的临时文件还是会写入HDFS,提交时添加配置
spark.hadoop.hadoop.tmp.dir=gs://<专用临时桶>/tmp/即可;注意spark.local.dir是节点本地磁盘的计算临时目录,和HDFS无关,保留本地磁盘配置即可,不需要改到GCS,否则会严重影响计算性能。 - 权限校验:提前确认服务账号对目标业务桶、临时桶都有完整的读写、列目录、删除对象权限,否则overwrite模式清理旧数据、写入提交标记文件时会报403权限错误。
- 提交示例:
spark-submit \ --class com.xxx.bi.etl.OrderDailyEtl \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --num-executors 20 \ --jars $HADOOP_HOME/share/hadoop/common/lib/gcs-connector-hadoop3-2.2.10.jar,$SPARK_HOME/jars/spark-avro_2.12-3.3.2.jar \ --conf spark.hadoop.google.cloud.auth.service.account.json.keyfile=/etc/hadoop/conf/gcs-prod-sa.json \ --conf spark.hadoop.hadoop.tmp.dir=gs://spark-tmp-prod/tmp/ \ etl-job-assembly-1.0.jar
Oozie调度注意事项
- ShareLib配置:提前把GCS connector、spark-avro的jar包上传到Oozie的ShareLib目录并执行更新,不要只把依赖放在单个工作流的lib目录下,否则Oozie自身读取GCS上的工作流定义、作业jar包时会报类找不到的错误。
- 路径全量替换:把workflow.xml、coordinator.xml中所有原来指向HDFS的路径全部替换为GCS路径,包括作业jar包路径、配置文件路径、输入输出数据路径,配置正确的情况下Oozie可以直接读取GCS上的资源,不需要再把资源同步到HDFS。
- 配置透传:如果没有在集群全局spark-defaults.conf、core-site.xml中配置GCS相关参数,需要在Spark action的
<spark-opts>标签中把GCS文件系统实现、密钥路径、临时路径这些配置透传给Spark作业,避免Executor端读不到配置。 - 全量弃用HDFS的额外配置:如果要完全移除HDFS的依赖,还需要把Oozie自身的系统依赖路径
oozie.service.WorkflowAppService.system.libpath从原来的HDFS路径改到GCS路径,同时给Oozie服务使用的服务账号配置对应路径的读写权限。
内容的提问来源于stack exchange,提问作者007chungking
相关产品推荐
相关产品推荐

