无法在EMR集群通过PySpark连接存储于S3的HBase
问题
在Amazon EMR emr-7.0.0集群(配套Spark 3.5.0、HBase 2.4.17)上,通过PySpark连接存储于S3的HBase时遇到报错,已通过HBase Shell确认目标表存在。
代码
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("test hbase").config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem").getOrCreate() spark.conf.set("hbase.zookeeper.quorum", "quorum_url") spark.conf.set("hbase.zookeeper.property.clientPort", "2181") spark.sparkContext.setSystemProperty("hbase.rootdir", "s3://emr-test/data/") spark.sparkContext.setSystemProperty("hbase.cluster.distributed", "true") spark.sparkContext.setSystemProperty("hbase.regionserver.global.memstore.upperLimit", "0.5") spark.read.format("org.apache.hadoop.hbase.spark").option("hbase.table", "example_table").load() print(hbase_tables) spark.stop()
报错信息
Traceback (most recent call last): File "<stdin>", line 1, in <module> File "/usr/lib/spark/python/pyspark/sql/readwriter.py", line 314, in load return self._df(self._jreader.load()) File "/usr/lib/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1322, in __call__ File "/usr/lib/spark/python/pyspark/errors/exceptions/captured.py", line 179, in deco return f(*a, **kw) File "/usr/lib/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/protocol.py", line 326, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o238.load. : org.apache.spark.SparkClassNotFoundException: [DATA_SOURCE_NOT_FOUND] Failed to find the data source: org.apache.hadoop.hbase.spark. Please find packages at `https://spark.apache.org/third-party-projects.html`.
解决方案
1. 核心原因
报错源于缺少HBase-Spark连接器JAR包,Spark默认不包含该组件,需手动引入适配版本的依赖。
2. 利用EMR预装JAR(推荐)
EMR集群通常已预装适配当前HBase版本的HBase-Spark连接器,路径为/usr/lib/hbase/lib/,文件名类似hbase-spark-2.4.17-amzn-*.jar。
修改spark-submit命令,通过--jars指定该JAR:
spark-submit --jars /usr/lib/hbase/lib/hbase-spark-2.4.17-amzn-1.jar test.py
若需批量引入HBase全量依赖,可直接指定lib目录:
spark-submit --jars $(echo /usr/lib/hbase/lib/*.jar | tr ' ' ',') test.py
3. 手动下载适配JAR
若集群本地无对应JAR,可下载适配Spark 3.5.0与HBase 2.4.17的连接器,Maven坐标为org.apache.hbase:hbase-spark:2.4.17。下载后上传至EMR主节点,再通过--jars参数指定路径即可。
4. 代码配置优化
- 无需手动设置
hbase.rootdir等集群级配置,直接让Spark读取HBase原生配置文件更可靠:
spark = SparkSession.builder.appName("test hbase") \ .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") \ .config("spark.driver.extraClassPath", "/usr/lib/hbase/conf:/usr/lib/hbase/lib/*") \ .config("spark.executor.extraClassPath", "/usr/lib/hbase/conf:/usr/lib/hbase/lib/*") \ .getOrCreate()
- 读取HBase表时必须指定列映射,否则无法解析数据,示例:
df = spark.read.format("org.apache.hadoop.hbase.spark") \ .option("hbase.table", "example_table") \ .option("hbase.columns.mapping", "cf1:col1,cf2:col2") # 替换为你的列族与列名 .load() df.show()
内容的提问来源于stack exchange,提问作者Asfandyar Abbasi
相关产品推荐
相关产品推荐

