Java实现ETL案例

wen java案例 1

目录导读

  1. ETL是什么?为什么Java是首选语言?
  2. 核心架构设计:一个轻量级ETL框架的模块拆解
  3. 实战案例:从CSV到MySQL的增量同步(含完整代码)
  4. 性能优化与异常处理的五个关键技巧
  5. 常见问题问答(FAQ)——解决你90%的踩坑点
  6. 下一步进阶路线图

ETL是什么?为什么Java是首选语言?

ETL(Extract-Transform-Load)是数据仓库建设的核心工序,分别对应抽取(Extract)转换(Transform)加载(Load),在现实业务中,它常被用于:将业务库(如Oracle)的数据同步到数仓(如Hive)、清洗日志文件、实时或准实时的数据集成。

Java实现ETL案例

为什么使用Java? 根据Tiobe 2024年度榜单,Java稳居前五,其优势体现在三方面:

  • 生态成熟:Apache Commons、Guava、以及Spring Batch等框架,让ETL开发不必重复造轮子。
  • 跨平台与健壮性:JVM的内存管理、异常处理机制,能扛住千万级数据的稳定跑批。
  • 并发模型:Java的ExecutorService和Fork/Join框架,能轻松实现并行抽取,而Python(GIL锁)或多线程编程更繁琐。

核心架构设计:一个轻量级ETL框架的模块拆解

一个标准Java ETL案例通常由四个模块构成,我们先看整体流程图(文字描述):

数据源(DB/文件/API) → [Extractor] → 中间数据集 → [Transformer] → 已清洗数据 → [Loader] → 目标库
                                    ↑                                            ↓
                              [Scheduler](可选:Quartz/Cron)        [Metrics Logger]

模块1:抽取器(Extractor)

  • 职责:负责连接数据源,分页读取或流式读取。
  • 要点:使用JDBC的fetchSize避免内存溢出;对于大文件,采用BufferedReader逐行读取。

模块2:转换器(Transformer)

  • 职责:清洗(去重、格式规整)、映射(字段改名)、计算(聚合、表达式)。
  • 要点:使用Map<String, Object>作为通用行模型,结合函数式接口Function进行链式处理,
    Function<Map<String,Object>, Map<String,Object>> cleanAge = row -> {
      int age = (int) row.get("age");
      row.put("age", age < 0 ? 0 : age);
      return row;
    };

模块3:加载器(Loader)

  • 职责:批量写入目标库,支持批量提交(addBatch)和幂等写入(先删后插或使用主键冲突更新)。

模块4:调度与监控

  • 使用@Scheduled(Spring)或Quartz触发任务,并用SLF4J记录每一批次的行数、耗时。

实战案例:从CSV到MySQL的增量同步(含完整代码)

场景需求:每天凌晨2点,将/data/orders_YYYYMMDD.csv文件(包含新订单)同步到MySQL的orders表,要求:

  • 如果订单ID已存在,则更新金额(amount)。
  • 处理时间要控制在10分钟以内(约50万行数据)。

步骤1:项目依赖(Maven)

<dependency>
    <groupId>com.zaxxer</groupId>
    <artifactId>HikariCP</artifactId>
    <version>4.0.3</version>
</dependency>
<dependency>
    <groupId>org.apache.commons</groupId>
    <artifactId>commons-csv</artifactId>
    <version>1.10.0</version>
</dependency>

步骤2:抽取器核心代码

public List<Order> extract(String filePath) throws IOException {
    List<Order> orders = new ArrayList<>();
    try (Reader reader = Files.newBufferedReader(Paths.get(filePath));
         CSVParser parser = new CSVParser(reader, CSVFormat.DEFAULT.withFirstRecordAsHeader())) {
        for (CSVRecord record : parser) {
            Order o = new Order();
            o.setId(Long.parseLong(record.get("order_id")));
            o.setAmount(Double.parseDouble(record.get("amount")));
            o.setStatus(record.get("status"));
            orders.add(o);
        }
    }
    return orders;
}

步骤3:转换与加载(合并了Transformer和Loader)

public void loadBatch(List<Order> orders) {
    String insertSql = "INSERT INTO orders (id, amount, status) VALUES (?,?,?) " +
                       "ON DUPLICATE KEY UPDATE amount = VALUES(amount), status = VALUES(status)";
    try (Connection conn = dataSource.getConnection();
         PreparedStatement ps = conn.prepareStatement(insertSql)) {
        conn.setAutoCommit(false);
        int count = 0;
        for (Order o : orders) {
            ps.setLong(1, o.getId());
            ps.setDouble(2, o.getAmount());
            ps.setString(3, o.getStatus());
            ps.addBatch();
            if (++count % 5000 == 0) {
                ps.executeBatch();
                conn.commit();
            }
        }
        ps.executeBatch();
        conn.commit();
    } catch (SQLException e) {
        // 记录错误批次,回滚并告警(这里省略日志框架)
    }
}

步骤4:主流程(Main方法简化)

public static void main(String[] args) {
    String date = LocalDate.now().minusDays(1).format(DateTimeFormatter.BASIC_ISO_DATE);
    String file = "/data/orders_" + date + ".csv";
    List<Order> data = extract(file);
    // 并发优化:使用并行流或线程池分片处理
    data.parallelStream().forEach(order -> transform(order));
    loadBatch(data);
}

性能优化与异常处理的五个关键技巧

批量读,批量写,绝不逐行写

  • 使用fetchSize(MySQL设为Integer.MIN_VALUE可强制流式读取),减少网络往返。

连接池与事务边界

  • 使用HikariCP(默认配置即可)管理连接,事务务必手动提交,避免每行自动提交导致性能雪崩。

并行化但控制线程数

  • 对于文件抽取,使用ForkJoinPool切分文件段;对于数据库写,建议线程数不超过CPU核数×2,否则锁竞争严重。

幂等性与断点续跑

  • 写操作加上ON DUPLICATE KEY或使用MERGE(H2/PostgreSQL)。
  • 记录processing_status表,每次跑批前检查上次成功位置,失败时从offset继续。

内存保护

  • 如果数据量超过可用堆内存,改用Streaming API(如FileReaderCharsetDecoder)代替List拼接,数据落地到临时文件后再转换。

常见问题问答(FAQ)

问题1:Java实现ETL案例时,如何防止内存溢出(OOM)? 答:分三个层面——数据层面用流式读取(如reader.lines()),容器层面设置-Xmx合理值(建议不超过物理内存的2/3),架构层面使用java.util.stream.Stream惰性求值,或者引入Apache Spark(Java API)做分布式处理。

问题2:当源数据有重复记录,怎么去重最优雅? 答:如果是全量同步,在Transformer中维护一个HashSet<Long>记录已见主键,但注意内存消耗,更高效法是利用数据库唯一索引,用INSERT IGNORE(MySQL)或ON CONFLICT DO NOTHING(PostgreSQL)让数据库去重,代码零改造。

问题3:数据转换时,如果日期格式不统一(如"2024/01/01"和"2024-01-01"),如何处理? 答:使用DateTimeFormatter的多格式解析器DateTimeFormatterBuilder.appendOptional(),或者写一个嵌套try-catch依次尝试LocalDate.parse的多个pattern,推荐后者代码更清晰。

问题4:Java实现ETL案例时,如何保证数据一致性(比如加载失败不产生脏数据)? 答:采用目标表临时表 + 两阶段提交策略:先写入orders_temp,全部成功后执行RENAME TABLE orders_temp TO orders(MySQL原子操作),或者用事务包裹整个批量加载,但注意大事务锁表风险,需权衡。

问题5:调度器选择Quartz好还是Spring Scheduler好? 答:Spring@Scheduled简单够用,但缺少分布式锁,若集群多节点部署,请用Quartz + JDBC JobStore,并配置@DisallowConcurrentExecution防重入,对于流式ETL(如Kafka),可考虑Kafka StreamsFlink,但已超出本题范围。


下一步进阶路线图

本文通过一个“CSV转MySQL”的Java实现ETL案例,带你走通了从抽取、转换到加载的完整闭环,但真正生产级的ETL还要考虑:

  • 元数据管理(记录每个字段的来源和口径)。
  • 数据质量规则(空值率、格式正则校验)。
  • 监控告警(Prometheus + Grafana埋点)。

建议你在此基础上尝试:

  1. 将转换逻辑抽成独立微服务,配合消息队列(RocketMQ/RabbitMQ)解耦。
  2. 学习Spring Batch框架,它提供了完善的ItemReader/Processor/Writer接口和重试机制,适合中小型批处理。
  3. 若数据量达到亿级,拥抱Apache Flink(Java实现)做流批一体。

数据管道永无止境,保持好奇,每次跑批的日志就是你最好的老师。

(本文完)

上一篇Canal案例

下一篇Kettle调用案例

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