Java大数据处理案例

wen java案例 2

本文目录导读:

Java大数据处理案例

  1. 环境准备
  2. 案例一:Hadoop MapReduce - 网站日志分析
  3. 案例二:Spark - 实时推荐系统
  4. 案例三:Flink - 实时流处理
  5. 案例四:批处理与数据库集成
  6. 使用示例
  7. 性能优化建议
  8. 注意事项

我将为您提供一个完整的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
    }
}

性能优化建议

  1. 内存优化:使用恰当的内存分配并启用序列化
  2. 并行度调整:根据数据量和集群规模调整并行度
  3. 数据压缩:使用Snappy等压缩算法减少网络传输
  4. 缓存策略:复用频繁访问的数据
  5. 计算优化:使用broadcast variables减少shuffle

注意事项

  1. 大数据处理需要合理设计数据分片
  2. 注意处理数据倾斜问题
  3. 做好故障恢复和容错处理
  4. 合理配置资源隔离策略

这个案例涵盖了Hadoop批处理、Spark内存计算、Flink流处理和数据库集成等主流大数据处理场景,可以作为学习和开发的参考。

上一篇Hadoop案例

下一篇Java车联网案例

抱歉,评论功能暂时关闭!