如何读取HBase导出到HDFS的数据并显示列名?
解决HBase导出数据Spark读取时显示列名的问题
方法一:导出时直接生成带列名的CSV
如果还没执行导出操作,直接修改hbase Export命令,让HBase输出包含列族+列名的CSV文件,后续Spark读取就能直接看到列名:
hbase org.apache.hadoop.hbase.mapreduce.Export \ <tablename> <outputdir> \ -D hbase.mapreduce.export.output.format=csv \ -D hbase.mapreduce.export.csv.separator=,
-D hbase.mapreduce.export.output.format=csv:指定导出格式为CSV,自动把列族和列名拼接成cf:col作为表头- 要是只需要导出特定列,在命令末尾加上列名即可,比如
cf1:col1 cf2:col2 - 导出完成后,直接用Spark的
spark.read.csv()读取文件,就能看到完整的列名和对应值
方法二:Spark中手动解析已导出的SequenceFile
如果已经用默认方式导出了SequenceFile,那就得放弃自带的转换器,手动解析Result对象来提取列名:
Python代码实现
from pyspark import SparkContext from pyspark.sql import Row from org.apache.hadoop.hbase.util import Bytes def parse_result(result): row_data = {"row_key": Bytes.toString(result.getRow())} # 遍历所有列族 for family in result.getMap().keySet(): family_bytes = Bytes.toBytes(family) # 遍历列族下的所有列 for qualifier in result.getFamilyMap(family_bytes).keySet(): col_name = f"{Bytes.toString(family_bytes)}:{Bytes.toString(qualifier)}" row_data[col_name] = Bytes.toString(result.getValue(family_bytes, qualifier)) return Row(**row_data) # 初始化SparkContext sc = SparkContext.getOrCreate() hadoop_conf = { "io.serializations": "org.apache.hadoop.io.serializer.WritableSerialization,org.apache.hadoop.hbase.mapreduce.ResultSerialization" } # 读取SequenceFile,获取原始Result对象 hbase_rdd = sc.newAPIHadoopFile( "<input_path>", 'org.apache.hadoop.mapreduce.lib.input.SequenceFileInputFormat', keyClass='org.apache.hadoop.hbase.io.ImmutableBytesWritable', valueClass='org.apache.hadoop.hbase.client.Result', conf=hadoop_conf ) # 解析Result,转换为带列名的Row parsed_rdd = hbase_rdd.map(lambda x: parse_result(x[1])) # 转成DataFrame,自动包含所有出现过的列名 df = parsed_rdd.toDF() # 查看前10条数据 df.show(10, truncate=False)
关键说明
- 不使用自带的
valueConverter,直接拿到Result原始对象,它包含了完整的列族、列名和值信息 - 通过
Result的API逐个提取列族和列限定符,拼接成可读的列名,再对应上列值 - 转换为DataFrame后,即使每行列名不同,Spark会自动收集所有出现过的列,缺失列的行会显示为
null
内容的提问来源于stack exchange,提问作者CompEng
相关产品推荐
相关产品推荐

