本文目录导读:

我将为您提供一个完整的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();
}
}
}
这个案例展示了:
- 管道设计模式:支持过滤、转换、输出等多个处理阶段
- Java 8+ 特性:Stream、Lambda、函数式接口
- 并行处理:支持多线程并行处理
- 异步处理:使用 CompletableFuture 实现异步操作
- 批处理:支持分批处理大量数据
- 错误处理:包含异常处理和回调机制
- 灵活性:可以轻松添加新的处理步骤
该案例可作为实际项目中数据处理管道的基础框架,可根据具体需求进行扩展和调整。