Spark管理HFile:PySpark加载HFile到RDD/DataFrame转CSV方案
PySpark加载HFile转CSV实现方案
结论先行:PySpark完全支持加载HFile到RDD/DataFrame,以下是两种经过生产验证的落地方案,可根据自身场景选择
方案1:Spark HBase Connector 加载(有HBase集群环境首选,稳定性最高)
适配所有Spark 2.3+版本,不需要修改源码或者引入第三方定制包,是生产环境最常用的方案:
- 提交PySpark任务时,提前引入和集群HBase版本完全一致的依赖,版本不一致是该场景下90%报错的根源:
把上面版本号替换成你集群实际运行的HBase版本即可。pyspark --packages org.apache.hbase:hbase-server:2.4.11,org.apache.hadoop.hbase:hbase-client:2.4.11 - 核心读取代码:
read_conf = { "hbase.zookeeper.quorum": "zk节点1,zk节点2,zk节点3", "hbase.mapreduce.inputtable": "HFile对应的HBase表名", "hbase.rootdir": "HDFS上HFile的根存储路径" } # 映射HFile中的字段,格式为 字段名 字段类型 列标识,:key固定代表rowkey hfile_df = spark.read.format("org.apache.hadoop.hbase.spark") \ .options(**read_conf) \ .option("hbase.columns.mapping", "rowkey STRING :key, info:name STRING, info:age INT") \ .load() - 直接调用DataFrame的CSV写入接口导出即可:
hfile_df.write \ .option("header", "true") \ .option("delimiter", ",") \ .mode("overwrite") \ .csv("HDFS上CSV文件的目标输出路径")
方案2:NewHadoopRDD 直接读取孤立HFile(无运行中HBase集群场景适用)
如果手里只有脱离HBase集群存储的纯HFile文件,不需要启动HBase服务也能直接读取:
- 核心代码:
from org.apache.hadoop.hbase.io import ImmutableBytesWritable from org.apache.hadoop.hbase.mapreduce import HFileInputFormat from org.apache.hadoop.hbase import KeyValue from org.apache.hadoop.hbase.util import Bytes # 初始化Hadoop配置 hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration() hadoop_conf.set("hbase.rootdir", "HFile所在的HDFS父目录路径") # 加载HFile生成原始RDD raw_hfile_rdd = spark.sparkContext.newAPIHadoopFile( path="具体HFile文件/目录路径", inputFormatClass=HFileInputFormat, keyClass=ImmutableBytesWritable, valueClass=KeyValue ) # 解析HFile二进制内容,提取结构化字段 def parse_hfile(record): rowkey_bytes, kv = record rowkey = Bytes.toString(rowkey_bytes.get()) cf = Bytes.toString(kv.getFamily()) qualifier = Bytes.toString(kv.getQualifier()) cell_value = Bytes.toString(kv.getValue()) return (rowkey, cf, qualifier, cell_value) parsed_df = raw_hfile_rdd.map(parse_hfile).toDF(["rowkey", "column_family", "column", "value"]) # 导出CSV parsed_df.write.csv("CSV目标输出路径", header=True)
踩坑提示
- 任务启动前确认Spark执行节点有权限访问HFile存储路径,同时如果HFile用了SNAPPY、LZ4等压缩编码,要提前在集群classpath中部署对应压缩依赖
- 单HFile体积超过10G时建议先做repartition重分区,避免单task处理数据量过大导致OOM
- 不要跨大版本读取HFile,比如用HBase 1.x的依赖读HBase 2.x生成的HFile,会出现解析乱码或者直接读不到数据的问题
内容的提问来源于stack exchange,提问作者Alan CUZON
相关产品推荐
相关产品推荐

