大学作业求助:基于AWS+Hadoop+Java SDK操作MySQL RDS的方案
Absolutely, your approach is totally feasible! This is a classic use case for combining Hadoop MapReduce's distributed data processing with AWS RDS MySQL's managed relational storage. Let me walk you through a step-by-step implementation with code examples and operational checklists.
Your design aligns perfectly with Hadoop's MapReduce paradigm: the Mapper parses raw data and groups records by A (the target table name), while the Reducer handles table creation and bulk data insertion into RDS MySQL. This distributed approach scales well for large datasets, and AWS EMR (Elastic MapReduce) integrates seamlessly with RDS as long as network access is configured correctly.
1. AWS Environment Setup
First, get your cloud resources configured:
- Create RDS MySQL Instance:
- Choose a suitable instance class (e.g.,
db.t3.microfor testing), set a master database name (e.g.,hadoop_rds_db), and note the endpoint, username, and password. - Update the RDS security group to allow inbound traffic on port 3306 from your EMR cluster's security group.
- Choose a suitable instance class (e.g.,
- Launch EMR Cluster:
- Select an EMR version with Hadoop (e.g., emr-6.10.0), ensure the cluster is in the same VPC as your RDS instance.
- Upload the MySQL JDBC driver (e.g.,
mysql-connector-java-8.0.30.jar) to/usr/lib/hadoop/lib/on the EMR master node (or package it into your MapReduce JAR).
2. MapReduce Code Implementation
Here's the Java code to implement your logic:
Mapper Class (Parses Raw Data)
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; public class DataParserMapper extends Mapper<Object, Text, Text, Text> { private Text tableKey = new Text(); private Text rowValue = new Text(); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { // Parse line format: (A, B, C) num1 num2 String line = value.toString().trim(); String innerPart = line.substring(line.indexOf('(')+1, line.indexOf(')')); String[] abc = innerPart.split(",\\s*"); String A = abc[0].trim(); String B = abc[1].trim(); String C = abc[2].trim(); String[] nums = line.substring(line.indexOf(')')+1).trim().split("\\s+"); String num1 = nums[0]; String num2 = nums[1]; // Output key = table name (A), value = pipe-separated row data tableKey.set(A); rowValue.set(B + "|" + C + "|" + num1 + "|" + num2); context.write(tableKey, rowValue); } }
Reducer Class (Creates Tables & Inserts Data)
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; import java.sql.SQLException; import java.util.Iterator; public class RDSDataReducer extends Reducer<Text, Text, Text, Text> { private Connection conn; @Override protected void setup(Context context) throws IOException, InterruptedException { // Fetch RDS config from Hadoop configuration String jdbcUrl = context.getConfiguration().get("rds.jdbc.url"); String username = context.getConfiguration().get("rds.username"); String password = context.getConfiguration().get("rds.password"); try { Class.forName("com.mysql.cj.jdbc.Driver"); conn = DriverManager.getConnection(jdbcUrl, username, password); conn.setAutoCommit(false); // Enable batch transactions } catch (ClassNotFoundException | SQLException e) { throw new IOException("Failed to connect to RDS MySQL", e); } } @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { String tableName = key.toString(); Iterator<Text> iter = values.iterator(); try { // Create table if it doesn't exist (B as primary key) String createTableSql = "CREATE TABLE IF NOT EXISTS `" + tableName + "` (" + "B VARCHAR(255) PRIMARY KEY, " + "C VARCHAR(255), " + "num1 INT, " + "num2 INT) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4"; try (PreparedStatement createStmt = conn.prepareStatement(createTableSql)) { createStmt.executeUpdate(); } // Batch insert/update data String insertSql = "INSERT INTO `" + tableName + "` (B, C, num1, num2) VALUES (?, ?, ?, ?) " + "ON DUPLICATE KEY UPDATE C=VALUES(C), num1=VALUES(num1), num2=VALUES(num2)"; try (PreparedStatement insertStmt = conn.prepareStatement(insertSql)) { int batchCount = 0; while (iter.hasNext()) { String[] rowData = iter.next().toString().split("\\|"); String B = rowData[0]; String C = rowData[1]; int num1 = Integer.parseInt(rowData[2]); int num2 = Integer.parseInt(rowData[3]); insertStmt.setString(1, B); insertStmt.setString(2, C); insertStmt.setInt(3, num1); insertStmt.setInt(4, num2); insertStmt.addBatch(); batchCount++; // Commit every 1000 records to optimize performance if (batchCount % 1000 == 0) { insertStmt.executeBatch(); conn.commit(); batchCount = 0; } } // Commit remaining records if (batchCount > 0) { insertStmt.executeBatch(); conn.commit(); } context.write(key, new Text("Successfully processed " + batchCount + " rows")); } } catch (SQLException e) { try { conn.rollback(); } catch (SQLException ex) { ex.printStackTrace(); } throw new IOException("Failed to process table " + tableName, e); } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { if (conn != null) { try { conn.close(); } catch (SQLException e) { e.printStackTrace(); } } } }
Driver Class (Configures & Runs the Job)
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.IOException; public class RDSDataIngestionDriver { public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException { Configuration conf = new Configuration(); // Set RDS connection parameters (replace with your values) conf.set("rds.jdbc.url", "jdbc:mysql://your-rds-endpoint:3306/hadoop_rds_db?useSSL=false&serverTimezone=UTC"); conf.set("rds.username", "your-rds-username"); conf.set("rds.password", "your-rds-password"); Job job = Job.getInstance(conf, "Hadoop-to-RDS-Data-Ingestion"); job.setJarByClass(RDSDataIngestionDriver.class); job.setMapperClass(DataParserMapper.class); job.setReducerClass(RDSDataReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // Set input/output paths (replace with your S3/HDFS paths) FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
- Package the Code: Compile the three classes into a JAR file (e.g.,
hadoop-rds-ingestion.jar). - Upload JAR to EMR: Transfer the JAR to the EMR master node or an S3 bucket.
- Run the MapReduce Job:
hadoop jar hadoop-rds-ingestion.jar com.yourpackage.RDSDataIngestionDriver s3://your-input-bucket/raw-data/ s3://your-output-bucket/job-logs/ - Verify Data in RDS: Connect to your RDS instance and check if tables are created with the correct data.
Once data is ingested, you can run standard SQL queries:
- Single Table Query:
SELECT B, C, num1 + num2 AS total FROM `A1` WHERE num1 > 50; - Join Tables:
SELECT a1.B, a1.C, a2.num2 FROM `A1` a1 INNER JOIN `A2` a2 ON a1.B = a2.B WHERE a1.num1 > 100;
- Table Name Sanitization: If
Acontains special characters, always wrap table names in backticks (`) to avoid SQL syntax errors. - Connection Pooling: For large clusters, use a connection pool (e.g., HikariCP) instead of creating a new JDBC connection per Reducer task to reduce overhead.
- Data Type Adjustments: Modify table column types (e.g., use
FLOATfor decimal numbers) based on your actual data. - Error Handling: Add more robust exception handling to handle data parsing failures (e.g., invalid
num1/num2values).
内容的提问来源于stack exchange,提问作者PeNpeL

