本文目录导读:

我将为您提供一个完整的Java大数据处理案例,包含多个实际场景的实现。
环境准备
<dependencies>
<!-- Hadoop -->
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.3.4</version>
</dependency>
<!-- Spark -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.12</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.12</artifactId>
<version>3.3.0</version>
</dependency>
<!-- Flink -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.15.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.12</artifactId>
<version>1.15.0</version>
</dependency>
<!-- 其他工具 -->
<dependency>
<groupId>commons-io</groupId>
<artifactId>commons-io</artifactId>
<version>2.11.0</version>
</dependency>
</dependencies>
案例一:Hadoop MapReduce - 网站日志分析
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import java.io.IOException;
import java.util.StringTokenizer;
import java.util.regex.Pattern;
public class LogAnalyzer {
// Mapper:解析日志并提取关键信息
public static class LogMapper extends Mapper<Object, Text, Text, IntWritable> {
private final static IntWritable one = new IntWritable(1);
private Text word = new Text();
private Pattern logPattern = Pattern.compile(
"(\\S+) (\\S+) (\\S+) \\[([^]]+)\\] \"([^\"]*)\" (\\d{3}) (\\d+)"
);
@Override
protected void map(Object key, Text value, Context context)
throws IOException, InterruptedException {
String line = value.toString();
java.util.regex.Matcher matcher = logPattern.matcher(line);
if (matcher.find()) {
String ip = matcher.group(1);
String timestamp = matcher.group(4);
String request = matcher.group(5);
String status = matcher.group(6);
// 提取访问的URL
String[] requestParts = request.split(" ");
if (requestParts.length > 1) {
String url = requestParts[1];
// 输出:URL + 状态码
String outputKey = url + " " + status;
word.set(outputKey);
context.write(word, one);
}
}
}
}
// Reducer:汇总统计
public static class LogReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
private IntWritable result = new IntWritable();
@Override
protected void reduce(Text key, Iterable<IntWritable> values,
Context context) throws IOException, InterruptedException {
int sum = 0;
for (IntWritable val : values) {
sum += val.get();
}
result.set(sum);
context.write(key, result);
}
}
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
Job job = Job.getInstance(conf, "log analyzer");
job.setJarByClass(LogAnalyzer.class);
job.setMapperClass(LogMapper.class);
job.setCombinerClass(LogReducer.class);
job.setReducerClass(LogReducer.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}
案例二:Spark - 实时推荐系统
import org.apache.spark.api.java.JavaPairRDD;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.sql.*;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructField;
import org.apache.spark.sql.types.StructType;
import scala.Tuple2;
import java.util.Arrays;
import java.util.List;
public class MovieRecommendation {
public static class Rating {
private int userId;
private int movieId;
private double rating;
public Rating(int userId, int movieId, double rating) {
this.userId = userId;
this.movieId = movieId;
this.rating = rating;
}
public int getUserId() { return userId; }
public int getMovieId() { return movieId; }
public double getRating() { return rating; }
}
public static void main(String[] args) {
SparkSession spark = SparkSession.builder()
.appName("MovieRecommendation")
.master("local[*]")
.getOrCreate();
// 1. 加载数据
Dataset<Row> ratingsDF = spark.read()
.option("header", "true")
.csv("ratings.csv");
// 2. 数据清洗和转换
ratingsDF = ratingsDF.select(
ratingsDF.col("userId").cast(DataTypes.IntegerType),
ratingsDF.col("movieId").cast(DataTypes.IntegerType),
ratingsDF.col("rating").cast(DataTypes.DoubleType)
);
ratingsDF.show(10);
ratingsDF.printSchema();
// 3. 基本统计分析
Dataset<Row> stats = ratingsDF
.groupBy("userId")
.agg(
functions.count("rating").alias("rating_count"),
functions.avg("rating").alias("avg_rating")
)
.orderBy(functions.desc("rating_count"));
stats.show(10);
// 4. 找出评分最高的电影
Dataset<Row> movieStats = ratingsDF
.groupBy("movieId")
.agg(
functions.count("rating").alias("num_ratings"),
functions.avg("rating").alias("avg_rating")
)
.filter("num_ratings > 50")
.orderBy(functions.desc("avg_rating"));
movieStats.show(10);
// 5. 简单的协同过滤推荐(基于物品的相似度)
JavaSparkContext jsc = new JavaSparkContext(spark.sparkContext());
JavaRDD<Rating> ratingsRDD = ratingsDF.javaRDD().map(row ->
new Rating(
row.getAs("userId"),
row.getAs("movieId"),
row.getAs("rating")
)
);
// 计算用户-物品矩阵
JavaPairRDD<Integer, Tuple2<Integer, Double>> userMoviePairs =
ratingsRDD.mapToPair(r ->
new Tuple2<>(r.getUserId(), new Tuple2<>(r.getMovieId(), r.getRating()))
);
// 计算物品之间的相似度
JavaPairRDD<Tuple2<Integer, Integer>, Double> moviePairs =
userMoviePairs.join(userMoviePairs)
.filter(pair -> !pair._2._1._1.equals(pair._2._2._1))
.mapToPair(pair -> {
int movie1 = pair._2._1._1;
int movie2 = pair._2._2._1;
double rating1 = pair._2._1._2;
double rating2 = pair._2._2._2;
Tuple2<Integer, Integer> movieKey =
movie1 < movie2 ?
new Tuple2<>(movie1, movie2) :
new Tuple2<>(movie2, movie1);
return new Tuple2<>(movieKey, Math.abs(rating1 - rating2));
})
.groupByKey()
.mapValues(ratings -> {
double sum = 0;
int count = 0;
for (Double rating : ratings) {
sum += rating;
count++;
}
return count > 0 ? sum / count : 0.0;
});
moviePairs.take(10).forEach(System.out::println);
// 6. 保存结果
movieStats.coalesce(1)
.write()
.mode(SaveMode.Overwrite)
.csv("movie_recommendations_result");
spark.stop();
}
}
案例三:Flink - 实时流处理
import org.apache.flink.api.common.functions.AggregateFunction;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.util.Collector;
public class RealTimeAnalytics {
public static void main(String[] args) throws Exception {
final ParameterTool params = ParameterTool.fromArgs(args);
// 创建执行环境
final StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.getConfig().setGlobalJobParameters(params);
// 1. 模拟数据流:用户点击流
DataStream<String> clickStream = env.addSource(new ClickSource());
// 2. 数据处理管道
DataStream<EventCount> counts = clickStream
.flatMap(new EventParser())
.keyBy(event -> event.getType())
.timeWindow(Time.minutes(1))
.aggregate(new EventAggregator());
// 3. 实时告警检测
DataStream<Alert> alerts = counts
.filter(count -> count.getCount() > 1000)
.map(count -> new Alert(
"HIGH_TRAFFIC",
count.getEventType(),
count.getCount(),
System.currentTimeMillis()
));
// 4. 输出结果
counts.print("Event Counts");
alerts.print("Alerts");
// 5. 实时用户行为分析
DataStream<UserBehavior> behaviorStream = clickStream
.flatMap(new UserBehaviorParser())
.keyBy(behavior -> behavior.getUserId())
.timeWindow(Time.minutes(5))
.aggregate(new UserBehaviorAggregator());
behaviorStream
.filter(behavior -> behavior.isSuspicious())
.print("Suspicious Users");
env.execute("Real-Time Analytics Pipeline");
}
// 数据源:模拟点击流
static class ClickSource implements org.apache.flink.streaming.api.functions.source.SourceFunction<String> {
private volatile boolean running = true;
@Override
public void run(SourceContext<String> ctx) throws Exception {
String[] events = {"click", "view", "purchase", "add_to_cart"};
int userId = 1;
while (running) {
String event = events[(int) (Math.random() * events.length)];
String data = String.format(
"%d,%s,%d,%d",
userId++,
event,
System.currentTimeMillis(),
(int) (Math.random() * 100)
);
ctx.collect(data);
Thread.sleep(100);
}
}
@Override
public void cancel() {
running = false;
}
}
// 事件解析
static class EventParser implements FlatMapFunction<String, Event> {
@Override
public void flatMap(String value, Collector<Event> out) {
String[] parts = value.split(",");
Event event = new Event(
Integer.parseInt(parts[0]),
parts[1],
Long.parseLong(parts[2]),
Integer.parseInt(parts[3])
);
out.collect(event);
}
}
// 聚合器
static class EventAggregator implements AggregateFunction<Event, EventCountAccumulator, EventCount> {
@Override
public EventCountAccumulator createAccumulator() {
return new EventCountAccumulator();
}
@Override
public EventCountAccumulator add(Event event, EventCountAccumulator acc) {
acc.add(event);
return acc;
}
@Override
public EventCount getResult(EventCountAccumulator acc) {
return new EventCount(acc.getType(), acc.getCount());
}
@Override
public EventCountAccumulator merge(EventCountAccumulator a, EventCountAccumulator b) {
a.merge(b);
return a;
}
}
// 简单的事件类
static class Event {
private int userId;
private String type;
private long timestamp;
private int value;
public Event(int userId, String type, long timestamp, int value) {
this.userId = userId;
this.type = type;
this.timestamp = timestamp;
this.value = value;
}
public String getType() { return type; }
public int getUserId() { return userId; }
}
// 事件计数类
static class EventCount {
private String eventType;
private long count;
public EventCount(String eventType, long count) {
this.eventType = eventType;
this.count = count;
}
public String getEventType() { return eventType; }
public long getCount() { return count; }
}
// 累加器类
static class EventCountAccumulator {
private String type;
private long count;
public void add(Event event) {
if (type == null) type = event.getType();
count++;
}
public void merge(EventCountAccumulator other) {
count += other.count;
}
public String getType() { return type; }
public long getCount() { return count; }
}
// 告警类
static class Alert {
private String type;
private String eventType;
private long count;
private long timestamp;
public Alert(String type, String eventType, long count, long timestamp) {
this.type = type;
this.eventType = eventType;
this.count = count;
this.timestamp = timestamp;
}
@Override
public String toString() {
return String.format(
"Alert{type='%s', eventType='%s', count=%d}",
type, eventType, count
);
}
}
// 用户行为类
static class UserBehavior {
private int userId;
private int eventCount;
private long totalValue;
public boolean isSuspicious() {
return eventCount > 100 && totalValue > 10000;
}
public int getUserId() { return userId; }
}
// 用户行为解析器
static class UserBehaviorParser implements FlatMapFunction<String, UserBehavior> {
@Override
public void flatMap(String value, Collector<UserBehavior> out) {
// 简化处理
}
}
// 用户行为聚合器
static class UserBehaviorAggregator implements AggregateFunction<UserBehavior, UserBehavior, UserBehavior> {
@Override
public UserBehavior createAccumulator() {
return new UserBehavior();
}
@Override
public UserBehavior add(UserBehavior value, UserBehavior acc) {
acc.eventCount++;
acc.totalValue += value.totalValue;
return acc;
}
@Override
public UserBehavior getResult(UserBehavior acc) {
return acc;
}
@Override
public UserBehavior merge(UserBehavior a, UserBehavior b) {
a.eventCount += b.eventCount;
a.totalValue += b.totalValue;
return a;
}
}
}
案例四:批处理与数据库集成
import java.sql.*;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
public class DatabaseBatchProcess {
private static final String JDBC_URL = "jdbc:mysql://localhost:3306/bigdata";
private static final String USERNAME = "root";
private static final String PASSWORD = "password";
public static void main(String[] args) throws Exception {
List<Record> records = readFromSource();
// 1. 批量插入
batchInsert(records);
// 2. 并行处理
parallelProcess(records);
// 3. 数据聚合
aggregateData();
}
// 批量插入
public static void batchInsert(List<Record> records) throws SQLException {
String sql = "INSERT INTO data_records (id, name, value, timestamp) VALUES (?, ?, ?, ?)";
try (Connection conn = DriverManager.getConnection(JDBC_URL, USERNAME, PASSWORD);
PreparedStatement ps = conn.prepareStatement(sql)) {
conn.setAutoCommit(false);
int batchSize = 1000;
int count = 0;
for (Record record : records) {
ps.setInt(1, record.getId());
ps.setString(2, record.getName());
ps.setDouble(3, record.getValue());
ps.setTimestamp(4, record.getTimestamp());
ps.addBatch();
if (++count % batchSize == 0) {
ps.executeBatch();
conn.commit();
}
}
ps.executeBatch();
conn.commit();
}
}
// 并行处理数据
public static void parallelProcess(List<Record> records) throws InterruptedException {
int numThreads = Runtime.getRuntime().availableProcessors();
ExecutorService executor = Executors.newFixedThreadPool(numThreads);
int chunkSize = records.size() / numThreads;
for (int i = 0; i < numThreads; i++) {
int start = i * chunkSize;
int end = (i == numThreads - 1) ? records.size() : (i + 1) * chunkSize;
List<Record> subList = records.subList(start, end);
executor.submit(() -> processChunk(subList));
}
executor.shutdown();
executor.awaitTermination(10, TimeUnit.MINUTES);
}
private static void processChunk(List<Record> records) {
for (Record record : records) {
// 处理每条记录
record.setProcessed(true);
}
}
// 数据聚合
public static void aggregateData() throws SQLException {
String sql = """
SELECT
DATE(timestamp) as date,
COUNT(*) as total_records,
AVG(value) as avg_value,
SUM(value) as total_value,
MIN(value) as min_value,
MAX(value) as max_value
FROM data_records
GROUP BY DATE(timestamp)
ORDER BY date DESC
""";
try (Connection conn = DriverManager.getConnection(JDBC_URL, USERNAME, PASSWORD);
Statement stmt = conn.createStatement();
ResultSet rs = stmt.executeQuery(sql)) {
while (rs.next()) {
System.out.printf(
"Date: %s, Count: %d, Avg: %.2f, Sum: %.2f%n",
rs.getDate("date"),
rs.getInt("total_records"),
rs.getDouble("avg_value"),
rs.getDouble("total_value")
);
}
}
}
// 数据记录类
static class Record {
private int id;
private String name;
private double value;
private Timestamp timestamp;
private boolean processed = false;
// getters and setters
public int getId() { return id; }
public void setId(int id) { this.id = id; }
public String getName() { return name; }
public void setName(String name) { this.name = name; }
public double getValue() { return value; }
public void setValue(double value) { this.value = value; }
public Timestamp getTimestamp() { return timestamp; }
public void setTimestamp(Timestamp timestamp) { this.timestamp = timestamp; }
public boolean isProcessed() { return processed; }
public void setProcessed(boolean processed) { this.processed = processed; }
}
// 从文件读取数据
private static List<Record> readFromSource() {
List<Record> records = new ArrayList<>();
// 模拟从文件或外部系统读取数据
return records;
}
}
使用示例
public class BigDataExample {
public static void main(String[] args) {
// 1. 运行 Hadoop 日志分析
// hadoop jar loganalyzer.jar LogAnalyzer /input/logs /output/analysis
// 2. 运行 Spark 推荐系统
// spark-submit --class MovieRecommendation movie-recommendation.jar
// 3. 运行 Flink 实时分析
// flink run -c RealTimeAnalytics real-time-analytics.jar
// 4. 数据库批处理
// java -cp .:lib/* DatabaseBatchProcess
}
}
性能优化建议
- 内存优化:使用恰当的内存分配并启用序列化
- 并行度调整:根据数据量和集群规模调整并行度
- 数据压缩:使用Snappy等压缩算法减少网络传输
- 缓存策略:复用频繁访问的数据
- 计算优化:使用broadcast variables减少shuffle
注意事项
- 大数据处理需要合理设计数据分片
- 注意处理数据倾斜问题
- 做好故障恢复和容错处理
- 合理配置资源隔离策略
这个案例涵盖了Hadoop批处理、Spark内存计算、Flink流处理和数据库集成等主流大数据处理场景,可以作为学习和开发的参考。