如何使用Spark将JPG、PDF及HDFS图片文件加载至HBase?
Absolutely! Spark is way more approachable and efficient than MapReduce for this task—perfect for someone getting started with the Hadoop ecosystem. It has a cleaner API, better in-memory performance, and easier debugging. Below’s a step-by-step guide tailored to your MOB-type HBase table (10MB threshold).
Prerequisites to Get Started
First, let’s make sure you’ve got the basics covered:
- A running Spark cluster compatible with your HBase version
- HBase config files (
hbase-site.xml,hbase-default.xml) copied to$SPARK_HOME/conf(so Spark can connect to HBase) - The HBase Spark connector JAR (match your HBase/Spark versions, e.g.,
hbase-spark-2.4.8.jar)
Quick Note on MOB Tables
Since your table uses MOB storage for objects over 10MB, you don’t need special code for MOB handling—HBase automatically stores files larger than the threshold as MOBs, while smaller files stay in the main table. Just write your data to the MOB-enabled column family, and HBase takes care of the rest.
Step-by-Step Implementation
1. Confirm Your HBase Table Setup
If you haven’t already, here’s how you’d create the MOB table via HBase Shell (for reference):
create 'file_store', {NAME => 'cf', IS_MOB => true, MOB_THRESHOLD => 10485760} # 10MB in bytes
We’ll use this table name (file_store) and column family (cf) in the code below.
2. Spark Code to Load HDFS Files into HBase
I’ll show both Scala (the most common for Spark-HBase integration) and Python (PySpark) snippets—pick whichever you’re more comfortable with.
Scala Implementation
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName} import org.apache.hadoop.hbase.client.Put import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapreduce.TableOutputFormat import org.apache.hadoop.hbase.util.Bytes import org.apache.spark.SparkContext import org.apache.spark.SparkConf object HdfsToHBaseMOB { def main(args: Array[String]): Unit = { // Initialize Spark app val conf = new SparkConf().setAppName("HdfsToHBaseMOB") val sc = new SparkContext(conf) // Load HBase config and set target table val hbaseConf = HBaseConfiguration.create() hbaseConf.set(TableOutputFormat.OUTPUT_TABLE, "file_store") sc.hadoopConfiguration.set(TableOutputFormat.OUTPUT_TABLE, "file_store") sc.hadoopConfiguration.set("mapreduce.job.outputformat.class", classOf[TableOutputFormat[ImmutableBytesWritable]].getName) // Path to your files in HDFS (supports wildcards like hdfs:///path/to/files/*) val hdfsFilePath = "hdfs:///user/your_username/documents/" // Read files from HDFS as binary data val fileRDD = sc.binaryFiles(hdfsFilePath) // Convert each file to a HBase Put object val hbaseRDD = fileRDD.map { case (filePath, contentStream) => // Use filename as row key (swap with a UUID if you have duplicate filenames) val rowKey = Bytes.toBytes(filePath.split("/").last) val put = new Put(rowKey) // Read file content into bytes val fileBytes = contentStream.toArray() // Add content to the MOB column family put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("content"), fileBytes) // Return the format HBase expects (new ImmutableBytesWritable, put) } // Write the data to HBase hbaseRDD.saveAsNewAPIHadoopDataset(sc.hadoopConfiguration) // Clean up Spark context sc.stop() } }
PySpark Implementation
You’ll need the happybase library installed on all Spark nodes first (pip install happybase):
from pyspark import SparkContext, SparkConf import happybase def main(): # Initialize Spark conf = SparkConf().setAppName("HdfsToHBaseMOB") sc = SparkContext(conf=conf) # HDFS path to your files hdfs_path = "hdfs:///user/your_username/documents/*" # Read files as binary data file_rdd = sc.binaryFiles(hdfs_path) # Connect to HBase (replace with your HBase master host) hbase_conn = happybase.Connection('hbase-master-hostname') table = hbase_conn.table('file_store') # Function to write each file to HBase def write_file_to_hbase(record): file_path, content = record # Use filename as row key (adjust if needed) row_key = file_path.split('/')[-1].encode('utf-8') # Write content to the MOB column family table.put(row_key, {b'cf:content': content}) # Execute the write operation on each file file_rdd.foreach(write_file_to_hbase) # Clean up connections hbase_conn.close() sc.stop() if __name__ == "__main__": main()
3. Submit the Spark Job
For Scala
Package your code into a JAR (using sbt or Maven), then submit with the required HBase JARs:
spark-submit \ --class HdfsToHBaseMOB \ --jars hbase-spark-<your-hbase-version>.jar,hbase-client-<your-hbase-version>.jar \ your-spark-jar-file.jar
For PySpark
Submit the script with the necessary dependencies:
spark-submit \ --packages com.hortonworks:shc-core:1.1.1-2.1-s_2.11 \ # Match your Spark/HBase version hdfs_to_hbase.py
Key Tips for Success
- Row Key Uniqueness: If you have duplicate filenames, use a unique identifier (like
UUID + filename) as the row key to avoid overwriting data. - Performance Tuning: Adjust Spark parameters like
--executor-memoryand--num-executorsto handle large file batches efficiently. - Validation: After writing, verify the data with HBase Shell:
get 'file_store', 'my-image.jpg', {COLUMNS => 'cf:content'}
内容的提问来源于stack exchange,提问作者Mutyam

