【发布时间】:2016-01-12 15:29:37
【问题描述】:
我们尝试测试以下用于访问 HBase 表(Spark-1.3.1、HBase-1.1.1、Hadoop-2.7.0)的示例代码:
import sys
from pyspark import SparkContext
if __name__ == "__main__":
if len(sys.argv) != 3:
print >> sys.stderr, """
Usage: hbase_inputformat <host> <table>
Run with example jar:
./bin/spark-submit --driver-class-path /path/to/example/jar \
/path/to/examples/hbase_inputformat.py <host> <table>
Assumes you have some data in HBase already, running on <host>, in <table>
"""
exit(-1)
host = sys.argv[1]
table = sys.argv[2]
sc = SparkContext(appName="HBaseInputFormat")
conf = {"hbase.zookeeper.quorum": host, "hbase.mapreduce.inputtable": table}
keyConv = "org.apache.spark.examples.pythonconverters.ImmutableBytesWritableToStringConverter"
valueConv = "org.apache.spark.examples.pythonconverters.HBaseResultToStringConverter"
hbase_rdd = sc.newAPIHadoopRDD(
"org.apache.hadoop.hbase.mapreduce.TableInputFormat",
"org.apache.hadoop.hbase.io.ImmutableBytesWritable",
"org.apache.hadoop.hbase.client.Result",
keyConverter=keyConv,
valueConverter=valueConv,
conf=conf)
output = hbase_rdd.collect()
for (k, v) in output:
print (k, v)
sc.stop()
我们收到以下错误:
15/10/14 12:46:24 INFO BlockManagerMaster:已注册的 BlockManager 回溯(最近一次通话最后): 文件“/opt/python/son.py”,第 30 行,在 配置=配置) newAPIHadoopRDD 中的文件“/usr/hdp/2.3.0.0-2557/spark/python/pyspark/context.py”,第 547 行 jconf,批处理大小) 调用中的文件“/usr/hdp/2.3.0.0-2557/spark/python/lib/py4j-0.8.2.1-src.zip/py4j/java_gateway.py”,第 538 行 文件“/usr/hdp/2.3.0.0-2557/spark/python/lib/py4j-0.8.2.1-src.zip/py4j/protocol.py”,第 300 行,在 get_return_value py4j.protocol.Py4JJavaError: 调用 z:org.apache.spark.api.python.PythonRDD.newAPIHadoopRDD 时出错。 : java.lang.ClassNotFoundException: org.apache.hadoop.hbase.io.ImmutableBytesWritable 在 java.net.URLClassLoader$1.run(URLClassLoader.java:366) 在 java.net.URLClassLoader$1.run(URLClassLoader.java:355) 在 java.security.AccessController.doPrivileged(本机方法) 在 java.net.URLClassLoader.findClass(URLClassLoader.java:354) 在 java.lang.ClassLoader.loadClass(ClassLoader.java:425) 在 java.lang.ClassLoader.loadClass(ClassLoader.java:358) 在 java.lang.Class.forName0(本机方法) 在 java.lang.Class.forName(Class.java:278) 在 org.apache.spark.util.Utils$.classForName(Utils.scala:157) 在 org.apache.spark.api.python.PythonRDD$.newAPIHadoopRDDFromClassNames(PythonRDD.scala:509) 在 org.apache.spark.api.python.PythonRDD$.newAPIHadoopRDD(PythonRDD.scala:494) 在 org.apache.spark.api.python.PythonRDD.newAPIHadoopRDD(PythonRDD.scala) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:606) 在 py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:231) 在 py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:379) 在 py4j.Gateway.invoke(Gateway.java:259) 在 py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:133) 在 py4j.commands.CallCommand.execute(CallCommand.java:79) 在 py4j.GatewayConnection.run(GatewayConnection.java:207) 在 java.lang.Thread.run(Thread.java:745)
高度赞赏任何见解。
【问题讨论】:
标签: hadoop apache-spark