如何在Dataproc集群的Jupyter Notebook中创建带Avro扩展的PySpark会话?
解决Spark 2.4.8读取Avro文件失败的问题
我们通过Concord自动创建启用PySpark的Dataproc集群,集群的PySpark Notebook使用Spark 2.4.8版本。默认情况下Spark没有Avro数据源扩展,无法读取.avro文件,尝试了如下配置但未生效:
PySpark会话配置
from pyspark.sql import SparkSession from pyspark import SparkConf import os # 设置Python路径 os.environ['PYSPARK_PYTHON'] = './PYENV1/pyenv1/bin/python3' os.environ['PYSPARK_DRIVER_PYTHON'] = './PYENV1/pyenv1/bin/python3' # 如果Notebook内核是PySpark,停止当前Spark应用 spark.sparkContext.stop() conf = SparkConf() # 添加配置让worker节点找到安装的包 conf.setAll([ ("spark.app.name", "Avro Testing"), \ ("spark.jars.packages","org.apache.spark:spark-avro_2.12:2.4.8"), \ # 其他注释的配置省略 ]) # conf.set("spark.jars","org.apache.spark:spark-avro_2.12:2.4.8")
读取数据代码
GCSPATH = "gs://gcs_bucket/file.avro" spark.read.format("avro").load(GCS_PATH)
错误信息
--------------------------------------------------------------------------- Py4JJavaError Traceback (most recent call last) /usr/lib/spark/python/pyspark/sql/utils.py in deco(*a, **kw) 62 try: ---> 63 return f(*a, **kw) 64 except py4j.protocol.Py4JJavaError as e: /usr/lib/spark/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name) 327 "An error occurred while calling {0}{1}{2}.\n". --> 328 format(target_id, ".", name), value) 329 else: Py4JJavaError: An error occurred while calling o622.load. : org.apache.spark.sql.AnalysisException: Failed to find data source: avro. Avro is built-in but external data source module since Spark 2.4. Please deploy the application as per the deployment section of "Apache Avro Data Source Guide".; at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:665) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:213) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:197) 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:750)
问题根源及解决办法
1. Scala版本不匹配(核心问题)
Spark 2.4.8默认基于Scala 2.11编译,你配置的spark-avro_2.12:2.4.8是Scala 2.12版本的包,两者不兼容,导致无法加载Avro数据源。
修正方式:把包名中的_2.12改成_2.11:
conf.set("spark.jars.packages", "org.apache.spark:spark-avro_2.11:2.4.8")
2. 未用新配置重新初始化SparkSession
你只创建了新的SparkConf,但没有用它重新构建SparkSession,原来的spark对象仍使用旧配置运行,根本没加载你指定的Avro包。
修正代码:在停止旧SparkContext后,用新的conf创建SparkSession:
# 停止旧的Spark应用后,添加这行代码 spark = SparkSession.builder.config(conf=conf).getOrCreate()
3. 集群层面的备选方案(全局支持Avro)
如果不想每次在Notebook里配置,可以在创建Dataproc集群时直接安装Avro扩展:
- 使用初始化动作:集群创建时添加
--initialization-actions gs://dataproc-initialization-actions/connectors/connectors.sh,并指定--metadata spark-avro-version=2.4.8 - 或者直接在集群创建命令中添加
--properties spark.jars.packages=org.apache.spark:spark-avro_2.11:2.4.8
内容的提问来源于stack exchange,提问作者Pratyay Sengupta
相关产品推荐
相关产品推荐

