Java数据处理管道案例

wen java案例 3

本文目录导读:

Java数据处理管道案例

  1. 基础数据处理管道设计
  2. 实际业务场景:用户数据处理
  3. 数据处理管道实际应用
  4. 增强版数据处理管道
  5. 使用示例

我将为您提供一个完整的Java数据处理管道案例,展示如何使用现代Java特性构建灵活、高效的数据处理流程。

基础数据处理管道设计

import java.util.*;
import java.util.function.*;
import java.util.stream.*;
import java.util.concurrent.*;
// 数据管道核心类
public class DataPipeline<T> {
    private final List<Function<T, T>> processors = new ArrayList<>();
    private final List<Predicate<T>> filters = new ArrayList<>();
    private final List<Consumer<T>> sinks = new ArrayList<>();
    // 添加处理器
    public DataPipeline<T> process(Function<T, T> processor) {
        processors.add(processor);
        return this;
    }
    // 添加过滤器
    public DataPipeline<T> filter(Predicate<T> predicate) {
        filters.add(predicate);
        return this;
    }
    // 添加输出目标
    public DataPipeline<T> sink(Consumer<T> sink) {
        sinks.add(sink);
        return this;
    }
    // 执行管道处理
    public void execute(Stream<T> source) {
        Stream<T> stream = source;
        // 应用过滤器
        for (Predicate<T> filter : filters) {
            stream = stream.filter(filter);
        }
        // 应用处理器
        for (Function<T, T> processor : processors) {
            stream = stream.map(processor);
        }
        // 输出结果
        List<T> results = stream.collect(Collectors.toList());
        results.forEach(item -> sinks.forEach(sink -> sink.accept(item)));
    }
    // 并行执行
    public void executeParallel(Stream<T> source, int threadPoolSize) {
        ForkJoinPool customThreadPool = new ForkJoinPool(threadPoolSize);
        try {
            customThreadPool.submit(() -> {
                Stream<T> stream = source.parallel();
                for (Predicate<T> filter : filters) {
                    stream = stream.filter(filter);
                }
                for (Function<T, T> processor : processors) {
                    stream = stream.map(processor);
                }
                List<T> results = stream.collect(Collectors.toList());
                results.forEach(item -> sinks.forEach(sink -> sink.accept(item)));
            }).get();
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            customThreadPool.shutdown();
        }
    }
}

实际业务场景:用户数据处理

import java.time.*;
import java.time.format.*;
// 用户数据类
class UserData {
    private String userId;
    private String name;
    private String email;
    private int age;
    private String city;
    private LocalDateTime registrationDate;
    private double totalPurchase;
    private int loginCount;
    public UserData(String userId, String name, String email, int age, 
                   String city, LocalDateTime registrationDate) {
        this.userId = userId;
        this.name = name;
        this.email = email;
        this.age = age;
        this.city = city;
        this.registrationDate = registrationDate;
        this.totalPurchase = 0.0;
        this.loginCount = 0;
    }
    // Getters and Setters
    public String getUserId() { return userId; }
    public void setUserId(String userId) { this.userId = userId; }
    public String getName() { return name; }
    public void setName(String name) { this.name = name; }
    public String getEmail() { return email; }
    public void setEmail(String email) { this.email = email; }
    public int getAge() { return age; }
    public void setAge(int age) { this.age = age; }
    public String getCity() { return city; }
    public void setCity(String city) { this.city = city; }
    public LocalDateTime getRegistrationDate() { return registrationDate; }
    public void setRegistrationDate(LocalDateTime registrationDate) { 
        this.registrationDate = registrationDate; 
    }
    public double getTotalPurchase() { return totalPurchase; }
    public void setTotalPurchase(double totalPurchase) { 
        this.totalPurchase = totalPurchase; 
    }
    public int getLoginCount() { return loginCount; }
    public void setLoginCount(int loginCount) { this.loginCount = loginCount; }
    @Override
    public String toString() {
        return String.format("UserData{id='%s', name='%s', email='%s', age=%d, city='%s'}",
                           userId, name, email, age, city);
    }
}
// 数据处理转换结果类
class ProcessedUser {
    private String userId;
    private String displayName;
    private String emailDomain;
    private String ageGroup;
    private String city;
    private boolean isActive;
    private double totalPurchase;
    public ProcessedUser(String userId, String displayName, String emailDomain,
                        String ageGroup, String city, boolean isActive, double totalPurchase) {
        this.userId = userId;
        this.displayName = displayName;
        this.emailDomain = emailDomain;
        this.ageGroup = ageGroup;
        this.city = city;
        this.isActive = isActive;
        this.totalPurchase = totalPurchase;
    }
    @Override
    public String toString() {
        return String.format("ProcessedUser{id='%s', name='%s', domain='%s', group='%s', city='%s', active=%b, purchase=%.2f}",
                           userId, displayName, emailDomain, ageGroup, city, isActive, totalPurchase);
    }
}

数据处理管道实际应用

import java.util.*;
import java.util.concurrent.*;
import java.util.stream.*;
public class PipelineDemo {
    // 生成测试数据
    private static List<UserData> generateTestData() {
        List<UserData> users = new ArrayList<>();
        Random random = new Random(42);
        String[] cities = {"北京", "上海", "广州", "深圳", "杭州", "成都"};
        String[] domains = {"gmail.com", "163.com", "qq.com", "outlook.com"};
        for (int i = 0; i < 1000; i++) {
            String userId = "USER_" + String.format("%04d", i);
            String name = "用户" + i;
            String email = "user" + i + "@" + domains[random.nextInt(domains.length)];
            int age = 18 + random.nextInt(50);
            String city = cities[random.nextInt(cities.length)];
            LocalDateTime regDate = LocalDateTime.now().minusDays(random.nextInt(365));
            UserData user = new UserData(userId, name, email, age, city, regDate);
            user.setTotalPurchase(100 + random.nextDouble() * 5000);
            user.setLoginCount(1 + random.nextInt(100));
            users.add(user);
        }
        return users;
    }
    // 数据处理管道示例
    public static void main(String[] args) {
        List<UserData> users = generateTestData();
        // 创建数据处理管道
        DataPipeline<UserData> pipeline = new DataPipeline<>();
        // 配置管道:过滤条件
        pipeline.filter(u -> u.getAge() >= 18 && u.getAge() <= 65)
               .filter(u -> u.getTotalPurchase() > 0)
               .filter(u -> u.getRegistrationDate().isBefore(LocalDateTime.now().minusMonths(1)));
        // 配置管道:数据处理
        pipeline.process(u -> {
            // 增强用户数据
            u.setLoginCount(u.getLoginCount() + 10);
            return u;
        }).process(u -> {
            // 添加数据转换
            return u;
        });
        // 配置输出:打印到控制台
        pipeline.sink(user -> {
            System.out.println("处理用户: " + user.getUserId());
        });
        // 配置输出:存储到数据库(模拟)
        List<UserData> processedUsers = new ArrayList<>();
        pipeline.sink(processedUsers::add);
        // 执行管道
        System.out.println("=== 单线程处理 ===");
        pipeline.execute(users.stream());
        // 并行处理示例
        System.out.println("\n=== 并行处理 ===");
        DataPipeline<UserData> parallelPipeline = new DataPipeline<>();
        parallelPipeline.filter(u -> u.getAge() >= 21)
                       .process(u -> {
                           u.setTotalPurchase(u.getTotalPurchase() * 1.1); // 10% 折扣
                           return u;
                       })
                       .sink(u -> {
                           System.out.println("VIP用户: " + u.getUserId() + 
                                            " 购买金额: " + u.getTotalPurchase());
                       });
        parallelPipeline.executeParallel(users.stream(), 4);
        // 高级数据处理:链式操作
        System.out.println("\n=== 统计分析 ===");
        analyzeUserData(users);
    }
    // 数据分析示例
    private static void analyzeUserData(List<UserData> users) {
        // 按城市分组统计
        Map<String, Long> cityCount = users.stream()
            .collect(Collectors.groupingBy(UserData::getCity, Collectors.counting()));
        System.out.println("城市分布:");
        cityCount.forEach((city, count) -> 
            System.out.printf("  %s: %d人%n", city, count));
        // 年龄分组
        Map<String, Long> ageGroups = users.stream()
            .collect(Collectors.groupingBy(
                u -> {
                    if (u.getAge() < 20) return "青少年";
                    else if (u.getAge() < 40) return "青年";
                    else if (u.getAge() < 60) return "中年";
                    else return "老年";
                },
                Collectors.counting()
            ));
        System.out.println("\n年龄分布:");
        ageGroups.forEach((group, count) -> 
            System.out.printf("  %s: %d人%n", group, count));
        // 消费统计
        double avgPurchase = users.stream()
            .mapToDouble(UserData::getTotalPurchase)
            .average()
            .orElse(0);
        double maxPurchase = users.stream()
            .mapToDouble(UserData::getTotalPurchase)
            .max()
            .orElse(0);
        System.out.println("\n消费统计:");
        System.out.printf("  平均消费: %.2f%n", avgPurchase);
        System.out.printf("  最高消费: %.2f%n", maxPurchase);
        // 活跃用户统计
        long activeUsers = users.stream()
            .filter(u -> u.getLoginCount() > 50)
            .count();
        System.out.println("\n活跃用户(登录>50次): " + activeUsers);
    }
    // 使用Java 8 Stream API的高级管道
    private static void advancedPipelineDemo() {
        List<UserData> users = generateTestData();
        // 链式 Stream 处理
        List<ProcessedUser> processed = users.stream()
            // 过滤
            .filter(u -> u.getAge() >= 18 && u.getAge() <= 60)
            .filter(u -> u.getTotalPurchase() > 1000)
            // 映射转换
            .map(u -> {
                String domain = u.getEmail().substring(u.getEmail().indexOf("@") + 1);
                String ageGroup = u.getAge() < 25 ? "年轻" : 
                                 (u.getAge() < 40 ? "青年" : "成年");
                boolean active = u.getLoginCount() > 10;
                return new ProcessedUser(
                    u.getUserId(),
                    u.getName(),
                    domain,
                    ageGroup,
                    u.getCity(),
                    active,
                    u.getTotalPurchase()
                );
            })
            // 排序
            .sorted(Comparator.comparingDouble(ProcessedUser::getTotalPurchase).reversed())
            // 限制数量
            .limit(50)
            // 收集结果
            .collect(Collectors.toList());
        // 输出统计
        System.out.println("高级管道处理结果:");
        processed.forEach(System.out::println);
    }
    // 批处理示例
    private static void batchProcessDemo() {
        List<UserData> users = generateTestData();
        // 批量处理工具
        BatchProcessor<UserData> batchProcessor = new BatchProcessor<>();
        batchProcessor.setBatchSize(5);
        // 使用批处理
        Iterable<List<UserData>> batches = batchProcessor.batch(users);
        for (List<UserData> batch : batches) {
            System.out.println("处理批次,大小: " + batch.size());
            // 每个批次的处理逻辑
            batch.parallelStream()
                .filter(u -> u.getTotalPurchase() > 0)
                .forEach(u -> {
                    // 模拟数据处理
                    double newPurchase = u.getTotalPurchase() * 1.2;
                    u.setTotalPurchase(newPurchase);
                });
        }
    }
    // 批处理辅助类
    static class BatchProcessor<T> {
        private int batchSize = 100;
        public void setBatchSize(int batchSize) {
            this.batchSize = batchSize;
        }
        public Iterable<List<T>> batch(List<T> items) {
            List<List<T>> batches = new ArrayList<>();
            for (int i = 0; i < items.size(); i += batchSize) {
                int end = Math.min(i + batchSize, items.size());
                batches.add(new ArrayList<>(items.subList(i, end)));
            }
            return batches;
        }
    }
}

增强版数据处理管道

// 支持类型转换的高级管道
public class AdvancedDataPipeline<I, O> {
    private final List<Processor<I, O>> processors = new ArrayList<>();
    @FunctionalInterface
    interface Processor<IN, OUT> {
        OUT process(IN input);
    }
    // 添加处理器
    public AdvancedDataPipeline<I, O> addProcessor(Function<I, O> processor) {
        processors.add(processor::apply);
        return this;
    }
    // 批量处理
    public List<O> processBatch(List<I> inputs) {
        return inputs.stream()
            .map(input -> {
                O result = null;
                for (Processor<I, O> processor : processors) {
                    result = processor.process(input);
                    input = (I) result; // 类型转换(注意类型安全)
                }
                return result;
            })
            .filter(Objects::nonNull)
            .collect(Collectors.toList());
    }
    // 带异常处理的管道
    public Optional<O> processWithErrorHandling(I input) {
        try {
            List<O> results = processBatch(Collections.singletonList(input));
            return results.isEmpty() ? Optional.empty() : Optional.of(results.get(0));
        } catch (Exception e) {
            System.err.println("处理错误: " + e.getMessage());
            return Optional.empty();
        }
    }
    // 延迟加载管道
    public Stream<O> lazyProcess(Stream<I> inputStream) {
        return inputStream.map(input -> {
            O result = null;
            for (Processor<I, O> processor : processors) {
                result = processor.process(input);
                input = (I) result;
            }
            return result;
        });
    }
    // 回调接口
    public interface PipelineListener<T> {
        void onData(T data);
        void onError(Exception e);
        void onComplete();
    }
    // 异步管道执行
    public CompletableFuture<List<O>> asyncProcess(List<I> inputs, PipelineListener<O> listener) {
        return CompletableFuture.supplyAsync(() -> {
            List<O> results = new ArrayList<>();
            for (I input : inputs) {
                try {
                    List<O> batchResults = processBatch(Collections.singletonList(input));
                    for (O result : batchResults) {
                        results.add(result);
                        if (listener != null) {
                            listener.onData(result);
                        }
                    }
                } catch (Exception e) {
                    if (listener != null) {
                        listener.onError(e);
                    }
                }
            }
            if (listener != null) {
                listener.onComplete();
            }
            return results;
        });
    }
}

使用示例

public class PipelineUsageExample {
    public static void main(String[] args) {
        // 创建输入数据
        List<String> rawData = Arrays.asList(
            "1,张三,25,北京,1000.50",
            "2,李四,30,上海,2000.00",
            "3,王五,22,广州,800.20",
            "4,赵六,35,深圳,3000.00"
        );
        // 创建高级管道
        AdvancedDataPipeline<String, Map<String, Object>> pipeline = 
            new AdvancedDataPipeline<>();
        // 添加处理步骤
        pipeline.addProcessor(line -> {
            String[] parts = line.split(",");
            Map<String, Object> map = new HashMap<>();
            map.put("id", Integer.parseInt(parts[0]));
            map.put("name", parts[1]);
            map.put("age", Integer.parseInt(parts[2]));
            map.put("city", parts[3]);
            map.put("purchaseAmount", Double.parseDouble(parts[4]));
            return map;
        });
        pipeline.addProcessor(data -> {
            // 添加年龄组信息
            int age = (int) data.get("age");
            String ageGroup = age < 25 ? "年轻" : (age < 35 ? "青年" : "中年");
            data.put("ageGroup", ageGroup);
            return data;
        });
        pipeline.addProcessor(data -> {
            // 添加VIP状态
            double amount = (double) data.get("purchaseAmount");
            data.put("isVIP", amount > 1500);
            data.put("discount", amount > 2000 ? 0.9 : 1.0);
            return data;
        });
        // 执行处理
        List<Map<String, Object>> processedData = pipeline.processBatch(rawData);
        // 输出结果
        System.out.println("处理结果:");
        processedData.forEach(System.out::println);
        // 异步处理
        AdvancedDataPipeline.PipelineListener<Map<String, Object>> listener = 
            new AdvancedDataPipeline.PipelineListener<Map<String, Object>>() {
                @Override
                public void onData(Map<String, Object> data) {
                    System.out.println("异步处理: " + data);
                }
                @Override
                public void onError(Exception e) {
                    System.err.println("异步错误: " + e.getMessage());
                }
                @Override
                public void onComplete() {
                    System.out.println("异步处理完成");
                }
            };
        // 启动异步处理
        pipeline.asyncProcess(rawData, listener)
            .whenComplete((result, error) -> {
                if (error != null) {
                    System.err.println("异步处理异常: " + error.getMessage());
                }
            });
        // 保持主线程运行以观察异步结果
        try {
            Thread.sleep(1000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

这个案例展示了:

  1. 管道设计模式:支持过滤、转换、输出等多个处理阶段
  2. Java 8+ 特性:Stream、Lambda、函数式接口
  3. 并行处理:支持多线程并行处理
  4. 异步处理:使用 CompletableFuture 实现异步操作
  5. 批处理:支持分批处理大量数据
  6. 错误处理:包含异常处理和回调机制
  7. 灵活性:可以轻松添加新的处理步骤

该案例可作为实际项目中数据处理管道的基础框架,可根据具体需求进行扩展和调整。

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