Java流批一体案例

wen java案例 3

Java流批一体案例:实时订单统计系统

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

Java流批一体案例

项目架构

┌─────────────────────────────────────────────────────────────┐
│                    业务场景:订单统计                         │
├─────────────────────────────────────────────────────────────┤
│  - 实时处理:实时计算每分钟的订单金额和数量                    │
│  - 批量处理:重算历史数据的订单统计                          │
│  - 统一输出:相同的数据格式和分析逻辑                        │
└─────────────────────────────────────────────────────────────┘

数据模型定义

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

核心优势

  1. 统一代码逻辑:流处理和批处理使用相同的解析、聚合和输出逻辑
  2. 降低维护成本:修改统计规则只需改一处代码
  3. 结果一致性:批处理结果和流处理结果可相互验证
  4. 灵活部署:同一套代码可根据需求切换运行模式
  5. 渐进式演进:先实现批处理,再平滑过渡到流处理

这个案例展示了如何用Flink实现真正的流批一体,大大简化了数据管道架构,提高了开发效率和系统可靠性。

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