Java流批一体案例:实时订单统计系统
流批一体是指使用同一套代码同时处理实时流数据和批量历史数据,下面我将通过一个订单统计系统的完整案例,展示如何使用Apache Flink实现流批一体。

项目架构
┌─────────────────────────────────────────────────────────────┐
│ 业务场景:订单统计 │
├─────────────────────────────────────────────────────────────┤
│ - 实时处理:实时计算每分钟的订单金额和数量 │
│ - 批量处理:重算历史数据的订单统计 │
│ - 统一输出:相同的数据格式和分析逻辑 │
└─────────────────────────────────────────────────────────────┘
数据模型定义
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonFormat;
import java.math.BigDecimal;
import java.time.LocalDateTime;
/**
* 订单事件 - 流和批共用同一数据模型
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderEvent {
private String orderId; // 订单ID
private String userId; // 用户ID
private String productId; // 商品ID
private BigDecimal amount; // 订单金额
@JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
private LocalDateTime orderTime; // 订单时间
private String status; // 订单状态:PAID/REFUND/CANCELLED
}
统计结果模型
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.math.BigDecimal;
/**
* 订单统计结果 - 流和批输出格式一致
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderStatistics {
private String windowStart; // 窗口开始时间
private String windowEnd; // 窗口结束时间
private Long orderCount; // 订单数量
private BigDecimal totalAmount; // 订单总金额
private BigDecimal avgAmount; // 平均订单金额
private String dataType; // 数据类型:REALTIME/BATCH
}
流批一体处理核心逻辑
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.connector.file.src.FileSource;
import org.apache.flink.connector.file.src.reader.TextLineInputFormat;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
import java.math.BigDecimal;
import java.time.Duration;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
/**
* 流批一体订单统计处理器
*/
public class UnifiedOrderStatisticsProcessor {
private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
/**
* 创建统一的处理管道 - 流和批共用
*/
public static SingleOutputStreamOperator<OrderStatistics> createProcessingPipeline(
DataStream<String> sourceStream,
String dataType) {
return sourceStream
// 1. 解析数据 - 流批共用解析逻辑
.map(new MapFunction<String, OrderEvent>() {
@Override
public OrderEvent map(String value) throws Exception {
return parseOrderEvent(value);
}
})
.filter(event -> event != null && "PAID".equals(event.getStatus()))
// 2. 设置事件时间和水位线 - 统一处理
.assignTimestampsAndWatermarks(
WatermarkStrategy.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((event, timestamp) ->
event.getOrderTime().toInstant(java.time.ZoneOffset.UTC).toEpochMilli())
)
// 3. 按照分钟进行窗口聚合 - 流批一致的窗口逻辑
.keyBy(new KeySelector<OrderEvent, String>() {
@Override
public String getKey(OrderEvent value) throws Exception {
return value.getProductId(); // 按商品分组
}
})
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
// 4. 聚合计算 - 统一的统计逻辑
.reduce(
new ReduceFunction<OrderEvent>() {
@Override
public OrderEvent reduce(OrderEvent event1, OrderEvent event2) throws Exception {
// 合并订单数据用于统计
OrderEvent merged = new OrderEvent();
merged.setOrderId(event1.getOrderId());
merged.setAmount(event1.getAmount().add(event2.getAmount()));
merged.setOrderTime(event1.getOrderTime());
// 使用特殊标记表示聚合结果
merged.setStatus("AGGREGATED");
return merged;
}
},
new ProcessWindowFunction<OrderEvent, OrderStatistics, String, TimeWindow>() {
@Override
public void process(String key, Context context,
Iterable<OrderEvent> elements,
Collector<OrderStatistics> out) throws Exception {
BigDecimal totalAmount = BigDecimal.ZERO;
Long orderCount = 0L;
for (OrderEvent event : elements) {
if ("AGGREGATED".equals(event.getStatus())) {
totalAmount = event.getAmount();
orderCount = 2L; // 简化处理,实际需要计数
} else {
totalAmount = totalAmount.add(event.getAmount());
orderCount++;
}
}
BigDecimal avgAmount = orderCount > 0
? totalAmount.divide(BigDecimal.valueOf(orderCount), 2, java.math.RoundingMode.HALF_UP)
: BigDecimal.ZERO;
out.collect(new OrderStatistics(
context.window().getStart() + "",
context.window().getEnd() + "",
orderCount,
totalAmount,
avgAmount,
dataType
));
}
}
);
}
/**
* 统一的解析逻辑 - 流和批共用
*/
private static OrderEvent parseOrderEvent(String line) {
try {
String[] fields = line.split(",");
if (fields.length < 6) return null;
OrderEvent event = new OrderEvent();
event.setOrderId(fields[0].trim());
event.setUserId(fields[1].trim());
event.setProductId(fields[2].trim());
event.setAmount(new BigDecimal(fields[3].trim()));
event.setOrderTime(LocalDateTime.parse(fields[4].trim(), FORMATTER));
event.setStatus(fields[5].trim());
return event;
} catch (Exception e) {
System.err.println("解析数据失败: " + line + ",错误: " + e.getMessage());
return null;
}
}
/**
* 批量处理 - 从文件读取历史数据
*/
public static DataStream<String> createBatchSource(
StreamExecutionEnvironment env,
String filePath) {
return env.readTextFile(filePath);
}
/**
* 流处理 - 从Kafka读取实时数据
*/
public static DataStream<String> createStreamSource(
StreamExecutionEnvironment env,
String kafkaTopic,
String bootstrapServers) {
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers(bootstrapServers)
.setTopics(kafkaTopic)
.setGroupId("order-statistics-group")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(
new org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema<String>() {
@Override
public void deserialize(
org.apache.kafka.clients.consumer.ConsumerRecord<byte[], byte[]> record,
Collector<String> out) {
out.collect(new String(record.value()));
}
@Override
public TypeInformation<String> getProducedType() {
return TypeInformation.of(String.class);
}
})
.build();
return env.fromSource(kafkaSource,
WatermarkStrategy.noWatermarks(),
"Kafka Source");
}
/**
* 主程序 - 演示流批一体
*/
public static void main(String[] args) throws Exception {
Configuration config = new Configuration();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
// 设置并行度
env.setParallelism(2);
// 1. 实时流处理
DataStream<String> streamSource = createStreamSource(
env,
"order-events",
"localhost:9092"
);
SingleOutputStreamOperator<OrderStatistics> realtimeResult =
createProcessingPipeline(streamSource, "REALTIME");
realtimeResult.print("实时统计结果: ");
// 2. 批量历史处理
DataStream<String> batchSource = createBatchSource(
env,
"/data/orders/history/2024-01-01.csv"
);
SingleOutputStreamOperator<OrderStatistics> batchResult =
createProcessingPipeline(batchSource, "BATCH");
batchResult.print("批量统计结果: ");
// 3. 统一输出到数据库或文件
// 实际项目中可以统一写入到MySQL、Elasticsearch或文件系统
env.execute("流批一体订单统计系统");
}
}
数据生成器(测试用)
import java.io.BufferedWriter;
import java.io.FileWriter;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Random;
/**
* 测试数据生成器
*/
public class DataGenerator {
private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
private static final String[] PRODUCTS = {"P001", "P002", "P003", "P004"};
private static final String[] USERS = {"U001", "U002", "U003", "U004", "U005"};
private static final Random RANDOM = new Random();
/**
* 生成历史批次数据文件
*/
public static void generateBatchData(String filePath, int count) throws Exception {
try (BufferedWriter writer = new BufferedWriter(new FileWriter(filePath))) {
for (int i = 0; i < count; i++) {
String orderId = "BATCH_ORDER_" + i;
String userId = USERS[RANDOM.nextInt(USERS.length)];
String productId = PRODUCTS[RANDOM.nextInt(PRODUCTS.length)];
BigDecimal amount = BigDecimal.valueOf(RANDOM.nextDouble() * 1000).setScale(2, java.math.RoundingMode.HALF_UP);
LocalDateTime time = LocalDateTime.now().minusDays(RANDOM.nextInt(30));
String status = RANDOM.nextDouble() > 0.2 ? "PAID" : "CANCELLED";
String line = String.format("%s,%s,%s,%s,%s,%s",
orderId, userId, productId, amount, time.format(FORMATTER), status);
writer.write(line);
writer.newLine();
}
}
}
/**
* 模拟实时流数据
*/
public static String generateStreamData() {
String orderId = "STREAM_ORDER_" + System.currentTimeMillis();
String userId = USERS[RANDOM.nextInt(USERS.length)];
String productId = PRODUCTS[RANDOM.nextInt(PRODUCTS.length)];
BigDecimal amount = BigDecimal.valueOf(RANDOM.nextDouble() * 1000).setScale(2, java.math.RoundingMode.HALF_UP);
LocalDateTime time = LocalDateTime.now();
String status = RANDOM.nextDouble() > 0.1 ? "PAID" : "REFUND";
return String.format("%s,%s,%s,%s,%s,%s",
orderId, userId, productId, amount, time.format(FORMATTER), status);
}
public static void main(String[] args) throws Exception {
// 生成测试数据
generateBatchData("/data/orders/history/2024-01-01.csv", 1000);
System.out.println("批量测试数据生成完成");
}
}
统一结果输出器
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.apache.flink.configuration.Configuration;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
/**
* 统一结果输出 - 流批结果写入同一张表
*/
public class UnifiedResultSink extends RichSinkFunction<OrderStatistics> {
private transient Connection connection;
private transient PreparedStatement statement;
@Override
public void open(Configuration parameters) throws Exception {
// 数据库连接配置
String url = "jdbc:mysql://localhost:3306/order_statistics";
String user = "root";
String password = "password";
connection = DriverManager.getConnection(url, user, password);
String sql = "INSERT INTO order_statistics " +
"(window_start, window_end, order_count, total_amount, avg_amount, data_type, create_time) " +
"VALUES (?, ?, ?, ?, ?, ?, NOW()) " +
"ON DUPLICATE KEY UPDATE " +
"order_count = VALUES(order_count), " +
"total_amount = VALUES(total_amount), " +
"avg_amount = VALUES(avg_amount), " +
"data_type = VALUES(data_type)";
statement = connection.prepareStatement(sql);
}
@Override
public void invoke(OrderStatistics value, Context context) throws Exception {
statement.setString(1, value.getWindowStart());
statement.setString(2, value.getWindowEnd());
statement.setLong(3, value.getOrderCount());
statement.setBigDecimal(4, value.getTotalAmount());
statement.setBigDecimal(5, value.getAvgAmount());
statement.setString(6, value.getDataType());
statement.executeUpdate();
}
@Override
public void close() throws Exception {
if (statement != null) statement.close();
if (connection != null) connection.close();
}
}
数据验证与对比示例
/**
* 流批结果对比验证
*/
public class StreamBatchComparison {
public static void main(String[] args) {
// 假设从数据库获取流和批的结果
List<OrderStatistics> realtimeResults = getRealtimeResults();
List<OrderStatistics> batchResults = getBatchResults();
// 按照时间窗口合并对比
Map<String, OrderStatistics> batchMap = batchResults.stream()
.collect(Collectors.toMap(
stat -> stat.getWindowStart() + "_" + stat.getWindowEnd(),
Function.identity()
));
for (OrderStatistics realtime : realtimeResults) {
String key = realtime.getWindowStart() + "_" + realtime.getWindowEnd();
OrderStatistics batch = batchMap.get(key);
if (batch != null) {
System.out.println("时间窗口: " + realtime.getWindowStart() + " - " + realtime.getWindowEnd());
System.out.println("实时统计: 数量=" + realtime.getOrderCount() +
", 金额=" + realtime.getTotalAmount());
System.out.println("批量统计: 数量=" + batch.getOrderCount() +
", 金额=" + batch.getTotalAmount());
System.out.println("差异: 数量=" +
Math.abs(realtime.getOrderCount() - batch.getOrderCount()) +
", 金额=" + realtime.getTotalAmount().subtract(batch.getTotalAmount()).abs());
System.out.println("---");
}
}
}
}
部署配置
# flink-conf.yaml 配置文件 # 流批一体配置 pipeline.name: UnifiedOrderStatistics # 检查点配置(批处理自动忽略) execution.checkpointing.interval: 60s execution.checkpointing.timeout: 10min # 状态后端 state.backend: rocksdb state.backend.incremental: true # 并行度 parallelism.default: 4 # 批处理优化 execution.runtime-mode: AUTOMATIC # 自动模式,根据数据源选择流或批
运行命令
# 提交Flink作业 - 流批一体模式 flink run -c com.example.UnifiedOrderStatisticsProcessor \ -D execution.runtime-mode=AUTOMATIC \ target/order-statistics-1.0.jar # 指定为批处理模式 flink run -c com.example.UnifiedOrderStatisticsProcessor \ -D execution.runtime-mode=BATCH \ target/order-statistics-1.0.jar # 指定为流处理模式 flink run -c com.example.UnifiedOrderStatisticsProcessor \ -D execution.runtime-mode=STREAMING \ target/order-statistics-1.0.jar
核心优势
- 统一代码逻辑:流处理和批处理使用相同的解析、聚合和输出逻辑
- 降低维护成本:修改统计规则只需改一处代码
- 结果一致性:批处理结果和流处理结果可相互验证
- 灵活部署:同一套代码可根据需求切换运行模式
- 渐进式演进:先实现批处理,再平滑过渡到流处理
这个案例展示了如何用Flink实现真正的流批一体,大大简化了数据管道架构,提高了开发效率和系统可靠性。