一个可在本地运行的入门级数据开发项目:从 CSV 文件抽取数据(Extract),做基础清洗转换(Transform),批量写入 PostgreSQL 数仓(Load)。


功能概述

  • Extract:从 data/input 下的 users.csv 和 orders.csv 读取数据
  • Transform:基础数据清洗
    • 去除多余空白
    • 邮箱校验与小写化
    • 金额转 BigDecimal
    • 时间格式解析
    • 去重(以 email 和 order_id 为唯一键)
    • 过滤脏数据(日志记录)
  • Load:使用 JDBC + HikariCP 批量 upsert 到 Postgres(ON CONFLICT)
  • Schema 管理:Flyway 自动迁移
  • 可观测性:SLF4J + Logback
  • 可运行环境:Docker Compose 一键启动 Postgres + Adminer

目录框架

data-pipeline-basic/
├── docker-compose.yml
├── pom.xml
├── data/
│   └── input/
│       ├── users.csv
│       └── orders.csv
├── src/
│   ├── main/
│   │   ├── java/com/example/datapipeline/
│   │   │   ├── Main.java
│   │   │   ├── config/Config.java
│   │   │   ├── db/DataSourceManager.java
│   │   │   ├── etl/CsvExtractor.java
│   │   │   ├── etl/JdbcLoader.java
│   │   │   ├── etl/TransformService.java
│   │   │   ├── model/RawOrderRecord.java
│   │   │   ├── model/RawUserRecord.java
│   │   │   ├── model/Order.java
│   │   │   ├── model/User.java
│   │   │   └── util/EmailValidator.java
│   │   └── resources/
│   │       ├── application.properties
│   │       ├── logback.xml
│   │       └── db/migration/V1__create_tables.sql
│   └── test/
│       └── java/com/example/datapipeline/etl/TransformServiceTest.java
└── README.md (本文即为说明)

快速开始

  1. 启动数据库(Postgres + Adminer)
  • 需安装 Docker 和 Docker Compose
  • 在项目根目录执行:
docker compose up -d
  • Postgres: localhost:5432,库 dwh,用户 dwh_user,密码 dwh_pass
  • Adminer: http://localhost:8080
    • System: PostgreSQL
    • Server: postgres
    • Username: dwh_user
    • Password: dwh_pass
    • Database: dwh
  1. 准备数据文件
  • 已提供样例:data/input/users.csv、data/input/orders.csv
  1. 运行
  • JDK 17+
  • Maven 3.8+
  • 先编译并打包可执行 jar:
mvn -q -DskipTests package
  • 运行:
java -jar target/data-pipeline-basic-1.0.0-shaded.jar
  1. 查看结果
  • 通过 Adminer 查询 users、orders 表数据
  • 程序日志输出在控制台(Logback)

可选:运行测试

mvn test

样例数据

data/input/users.csv

user_id,email,full_name,signup_date
1, Alice@example.com ,Alice Zhang,2023-01-10
2,bob@example.com,Bob Li,2023-03-22
3,invalid-email,Bad User,2023-05-01
4,carol@example.com,  Carol Wang  ,2023-06-15
2,bob@example.com,Bob Li Jr,2023-03-22

data/input/orders.csv

order_id,user_email,amount,currency,created_at,status
1001,alice@example.com,19.99,USD,2023-12-31T23:59:59Z,PAID
1002, bob@example.com ,29.50,USD,2024-01-01T08:00:00+08:00,PAID
1003,carol@example.com,12.00,USD,2024-01-02T10:30:00Z,CANCELLED
1004,unknown@example.com,5.00,USD,2024-01-03T10:30:00Z,PAID
1002,bob@example.com,30.00,USD,2024-01-05T09:00:00Z,REFUND
badrow,bademail,xx,USD,not-a-time,PENDING

说明:

  • users.csv 中 invalid-email 会被过滤,Bob 会被 upsert 成 “Bob Li Jr”
  • orders.csv 中 unknown@example.com 会被过滤(找不到用户邮箱),1002 为 upsert(状态更新)

全部代码

pom.xml

<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
         https://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>

  <groupId>com.example</groupId>
  <artifactId>data-pipeline-basic</artifactId>
  <version>1.0.0</version>

  <properties>
    <maven.compiler.source>17</maven.compiler.source>
    <maven.compiler.target>17</maven.compiler.target>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <slf4j.version>2.0.12</slf4j.version>
    <logback.version>1.4.14</logback.version>
    <junit.version>5.10.2</junit.version>
  </properties>

  <dependencies>
    <dependency>
      <groupId>org.slf4j</groupId>
      <artifactId>slf4j-api</artifactId>
      <version>${slf4j.version}</version>
    </dependency>
    <dependency>
      <groupId>ch.qos.logback</groupId>
      <artifactId>logback-classic</artifactId>
      <version>${logback.version}</version>
    </dependency>

    <dependency>
      <groupId>com.opencsv</groupId>
      <artifactId>opencsv</artifactId>
      <version>5.9</version>
    </dependency>

    <dependency>
      <groupId>com.zaxxer</groupId>
      <artifactId>HikariCP</artifactId>
      <version>5.1.0</version>
    </dependency>
    <dependency>
      <groupId>org.postgresql</groupId>
      <artifactId>postgresql</artifactId>
      <version>42.7.3</version>
    </dependency>

    <dependency>
      <groupId>org.flywaydb</groupId>
      <artifactId>flyway-core</artifactId>
      <version>10.15.0</version>
    </dependency>

    <dependency>
      <groupId>org.junit.jupiter</groupId>
      <artifactId>junit-jupiter</artifactId>
      <version>${junit.version}</version>
      <scope>test</scope>
    </dependency>
  </dependencies>

  <build>
    <plugins>
      <plugin>
        <artifactId>maven-surefire-plugin</artifactId>
        <version>3.2.5</version>
        <configuration>
          <useModulePath>false</useModulePath>
        </configuration>
      </plugin>

      <plugin>
        <groupId>org.apache.maven.plugins</groupId>
        <artifactId>maven-shade-plugin</artifactId>
        <version>3.5.0</version>
        <executions>
          <execution>
            <phase>package</phase>
            <goals><goal>shade</goal></goals>
            <configuration>
              <createDependencyReducedPom>false</createDependencyReducedPom>
              <transformers>
                <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                  <mainClass>com.example.datapipeline.Main</mainClass>
                </transformer>
              </transformers>
            </configuration>
          </execution>
        </executions>
      </plugin>

    </plugins>
  </build>
</project>

docker-compose.yml

version: "3.9"
services:
  postgres:
    image: postgres:16
    container_name: dwh-postgres
    environment:
      POSTGRES_DB: dwh
      POSTGRES_USER: dwh_user
      POSTGRES_PASSWORD: dwh_pass
    ports:
      - "5432:5432"
    volumes:
      - pgdata:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U dwh_user -d dwh"]
      interval: 3s
      timeout: 3s
      retries: 10

  adminer:
    image: adminer:4
    container_name: dwh-adminer
    ports:
      - "8080:8080"
    depends_on:
      - postgres

volumes:
  pgdata:

src/main/resources/application.properties

# 覆盖优先级:系统属性 > 环境变量 > 本文件
db.url=jdbc:postgresql://localhost:5432/dwh
db.username=dwh_user
db.password=dwh_pass

input.dir=./data/input

src/main/resources/logback.xml

<configuration>
  <appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
    <encoder>
      <pattern>%d{HH:mm:ss.SSS} %-5level [%thread] %logger{36} - %msg%n</pattern>
    </encoder>
  </appender>
  <root level="INFO">
    <appender-ref ref="STDOUT" />
  </root>
</configuration>

src/main/resources/db/migration/V1__create_tables.sql

CREATE TABLE IF NOT EXISTS users (
    user_id      BIGINT,
    email        TEXT NOT NULL UNIQUE,
    full_name    TEXT,
    signup_date  DATE,
    created_at   TIMESTAMPTZ DEFAULT NOW(),
    updated_at   TIMESTAMPTZ DEFAULT NOW()
);

CREATE TABLE IF NOT EXISTS orders (
    order_id     BIGINT PRIMARY KEY,
    user_email   TEXT NOT NULL,
    amount       NUMERIC(18,2) NOT NULL,
    currency     TEXT NOT NULL,
    created_at   TIMESTAMPTZ NOT NULL,
    status       TEXT NOT NULL,
    inserted_at  TIMESTAMPTZ DEFAULT NOW(),
    updated_at   TIMESTAMPTZ DEFAULT NOW(),
    CONSTRAINT fk_orders_user_email
      FOREIGN KEY (user_email) REFERENCES users(email)
);

-- 更新触发器(可选,这里简化用应用层维护 updated_at)

src/main/java/com/example/datapipeline/Main.java

package com.example.datapipeline;

import com.example.datapipeline.config.Config;
import com.example.datapipeline.db.DataSourceManager;
import com.example.datapipeline.etl.CsvExtractor;
import com.example.datapipeline.etl.JdbcLoader;
import com.example.datapipeline.etl.TransformService;
import com.example.datapipeline.model.Order;
import com.example.datapipeline.model.RawOrderRecord;
import com.example.datapipeline.model.RawUserRecord;
import com.example.datapipeline.model.User;
import com.zaxxer.hikari.HikariDataSource;
import org.flywaydb.core.Flyway;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.time.Duration;
import java.time.Instant;
import java.util.List;

public class Main {
    private static final Logger log = LoggerFactory.getLogger(Main.class);

    public static void main(String[] args) {
        Instant start = Instant.now();
        try {
            Config config = Config.load();
            log.info("Starting ETL with input.dir={} db.url={}", config.getInputDir(), config.getDbUrl());

            // Flyway
            Flyway flyway = Flyway.configure()
                    .dataSource(config.getDbUrl(), config.getDbUsername(), config.getDbPassword())
                    .load();
            flyway.migrate();
            log.info("Flyway migration completed");

            try (HikariDataSource ds = DataSourceManager.create(config)) {
                CsvExtractor extractor = new CsvExtractor(config.getInputDir());
                List<RawUserRecord> rawUsers = extractor.extractUsers();
                List<RawOrderRecord> rawOrders = extractor.extractOrders();

                TransformService transform = new TransformService();
                List<User> users = transform.transformUsers(rawUsers);
                List<Order> orders = transform.transformOrders(rawOrders, users);

                JdbcLoader loader = new JdbcLoader(ds);
                int upsertUsers = loader.upsertUsers(users);
                int upsertOrders = loader.upsertOrders(orders);

                log.info("ETL done. users: raw={} -> valid={} -> upserted={}, orders: raw={} -> valid={} -> upserted={}",
                        rawUsers.size(), users.size(), upsertUsers,
                        rawOrders.size(), orders.size(), upsertOrders);
            }

            Duration d = Duration.between(start, Instant.now());
            log.info("Total elapsed: {} ms", d.toMillis());
        } catch (Exception e) {
            log.error("ETL failed", e);
            System.exit(1);
        }
    }
}

src/main/java/com/example/datapipeline/config/Config.java

package com.example.datapipeline.config;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.InputStream;
import java.util.Objects;
import java.util.Properties;

public class Config {
    private static final Logger log = LoggerFactory.getLogger(Config.class);

    private final String dbUrl;
    private final String dbUsername;
    private final String dbPassword;
    private final String inputDir;

    private Config(String dbUrl, String dbUsername, String dbPassword, String inputDir) {
        this.dbUrl = Objects.requireNonNull(dbUrl);
        this.dbUsername = Objects.requireNonNull(dbUsername);
        this.dbPassword = Objects.requireNonNull(dbPassword);
        this.inputDir = Objects.requireNonNull(inputDir);
    }

    public static Config load() {
        Properties p = new Properties();
        try (InputStream in = Config.class.getClassLoader().getResourceAsStream("application.properties")) {
            if (in != null) {
                p.load(in);
            }
        } catch (Exception e) {
            log.warn("Failed to load application.properties, using defaults/env", e);
        }

        String dbUrl = override("DB_URL", "db.url", p, "jdbc:postgresql://localhost:5432/dwh");
        String dbUser = override("DB_USERNAME", "db.username", p, "dwh_user");
        String dbPass = override("DB_PASSWORD", "db.password", p, "dwh_pass");
        String inputDir = override("INPUT_DIR", "input.dir", p, "./data/input");

        return new Config(dbUrl, dbUser, dbPass, inputDir);
    }

    private static String override(String env, String propKey, Properties p, String def) {
        String sys = System.getProperty(propKey);
        if (sys != null && !sys.isBlank()) return sys;
        String e = System.getenv(env);
        if (e != null && !e.isBlank()) return e;
        String v = p.getProperty(propKey);
        if (v != null && !v.isBlank()) return v;
        return def;
    }

    public String getDbUrl() { return dbUrl; }
    public String getDbUsername() { return dbUsername; }
    public String getDbPassword() { return dbPassword; }
    public String getInputDir() { return inputDir; }
}

src/main/java/com/example/datapipeline/db/DataSourceManager.java

package com.example.datapipeline.db;

import com.example.datapipeline.config.Config;
import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;

public class DataSourceManager {

    public static HikariDataSource create(Config cfg) {
        HikariConfig hc = new HikariConfig();
        hc.setJdbcUrl(cfg.getDbUrl());
        hc.setUsername(cfg.getDbUsername());
        hc.setPassword(cfg.getDbPassword());
        hc.setMaximumPoolSize(4);
        hc.setMinimumIdle(1);
        hc.setPoolName("etl-pool");
        hc.setAutoCommit(false);
        return new HikariDataSource(hc);
    }
}

src/main/java/com/example/datapipeline/etl/CsvExtractor.java

package com.example.datapipeline.etl;

import com.example.datapipeline.model.RawOrderRecord;
import com.example.datapipeline.model.RawUserRecord;
import com.opencsv.CSVReaderHeaderAware;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.FileReader;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;

public class CsvExtractor {
    private static final Logger log = LoggerFactory.getLogger(CsvExtractor.class);

    private final Path inputDir;

    public CsvExtractor(String inputDir) {
        this.inputDir = Path.of(inputDir);
    }

    public List<RawUserRecord> extractUsers() {
        Path file = inputDir.resolve("users.csv");
        List<RawUserRecord> out = new ArrayList<>();
        try (CSVReaderHeaderAware reader = new CSVReaderHeaderAware(new FileReader(file.toFile()))) {
            Map<String, String> row;
            while ((row = reader.readMap()) != null) {
                RawUserRecord r = new RawUserRecord(
                        row.getOrDefault("user_id", ""),
                        row.getOrDefault("email", ""),
                        row.getOrDefault("full_name", ""),
                        row.getOrDefault("signup_date", "")
                );
                out.add(r);
            }
            log.info("Extracted {} user rows from {}", out.size(), file);
        } catch (Exception e) {
            throw new RuntimeException("Failed to read users.csv from " + file, e);
        }
        return out;
    }

    public List<RawOrderRecord> extractOrders() {
        Path file = inputDir.resolve("orders.csv");
        List<RawOrderRecord> out = new ArrayList<>();
        try (CSVReaderHeaderAware reader = new CSVReaderHeaderAware(new FileReader(file.toFile()))) {
            Map<String, String> row;
            while ((row = reader.readMap()) != null) {
                RawOrderRecord r = new RawOrderRecord(
                        row.getOrDefault("order_id", ""),
                        row.getOrDefault("user_email", ""),
                        row.getOrDefault("amount", ""),
                        row.getOrDefault("currency", ""),
                        row.getOrDefault("created_at", ""),
                        row.getOrDefault("status", "")
                );
                out.add(r);
            }
            log.info("Extracted {} order rows from {}", out.size(), file);
        } catch (Exception e) {
            throw new RuntimeException("Failed to read orders.csv from " + file, e);
        }
        return out;
    }
}

src/main/java/com/example/datapipeline/etl/TransformService.java

package com.example.datapipeline.etl;

import com.example.datapipeline.model.*;
import com.example.datapipeline.util.EmailValidator;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.*;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.stream.Collectors;

public class TransformService {
    private static final Logger log = LoggerFactory.getLogger(TransformService.class);

    private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ISO_LOCAL_DATE;

    public List<User> transformUsers(List<RawUserRecord> raw) {
        Map<String, User> dedupByEmail = new LinkedHashMap<>();
        int dropped = 0;

        for (RawUserRecord r : raw) {
            String email = safeTrim(r.email()).toLowerCase(Locale.ROOT);
            String fullName = safeTrim(r.fullName());
            String signupDateStr = safeTrim(r.signupDate());
            String userIdStr = safeTrim(r.userId());

            if (!EmailValidator.isValid(email)) {
                dropped++;
                log.warn("Drop user row due to invalid email: {}", r);
                continue;
            }

            LocalDate signupDate = null;
            if (!signupDateStr.isBlank()) {
                try {
                    signupDate = LocalDate.parse(signupDateStr, DATE_FMT);
                } catch (Exception e) {
                    log.warn("Invalid signup_date '{}', set null. row={}", signupDateStr, r);
                }
            }

            Long userId = null;
            try {
                if (!userIdStr.isBlank()) {
                    userId = Long.parseLong(userIdStr);
                }
            } catch (Exception e) {
                log.warn("Invalid user_id '{}', set null. row={}", userIdStr, r);
            }

            User u = new User(userId, email, fullName, signupDate);
            // 保留最后一次出现的记录(示例:更新名字)
            dedupByEmail.put(email, u);
        }

        if (dropped > 0) {
            log.info("Users transform: dropped {} invalid rows", dropped);
        }
        return new ArrayList<>(dedupByEmail.values());
    }

    public List<Order> transformOrders(List<RawOrderRecord> rawOrders, List<User> validUsers) {
        Set<String> knownEmails = validUsers.stream().map(User::email).collect(Collectors.toSet());
        Map<Long, Order> dedupByOrderId = new LinkedHashMap<>();
        int dropped = 0;

        for (RawOrderRecord r : rawOrders) {
            String orderIdStr = safeTrim(r.orderId());
            String email = safeTrim(r.userEmail()).toLowerCase(Locale.ROOT);
            String amountStr = safeTrim(r.amount());
            String currency = safeTrim(r.currency()).toUpperCase(Locale.ROOT);
            String createdAtStr = safeTrim(r.createdAt());
            String status = safeTrim(r.status()).toUpperCase(Locale.ROOT);

            Long orderId;
            try {
                orderId = Long.parseLong(orderIdStr);
            } catch (Exception e) {
                dropped++;
                log.warn("Drop order row due to invalid order_id: {}", r);
                continue;
            }

            if (!EmailValidator.isValid(email) || !knownEmails.contains(email)) {
                dropped++;
                log.warn("Drop order row due to unknown/invalid email: {}", r);
                continue;
            }

            BigDecimal amount;
            try {
                amount = new BigDecimal(amountStr).setScale(2, RoundingMode.HALF_UP);
            } catch (Exception e) {
                dropped++;
                log.warn("Drop order row due to invalid amount: {}", r);
                continue;
            }
            if (currency.isBlank()) currency = "USD";

            OffsetDateTime createdAt;
            try {
                // 支持 Z 或者含偏移量的时间
                createdAt = OffsetDateTime.parse(createdAtStr);
            } catch (Exception e) {
                try {
                    // 尝试无偏移量当作 UTC
                    LocalDateTime ldt = LocalDateTime.parse(createdAtStr);
                    createdAt = ldt.atOffset(ZoneOffset.UTC);
                } catch (Exception ex) {
                    dropped++;
                    log.warn("Drop order row due to invalid created_at: {}", r);
                    continue;
                }
            }
            if (status.isBlank()) status = "PENDING";

            Order o = new Order(orderId, email, amount, currency, createdAt, status);
            // 保留最后一次(示例:更新订单状态/金额)
            dedupByOrderId.put(orderId, o);
        }

        if (dropped > 0) {
            log.info("Orders transform: dropped {} invalid rows", dropped);
        }
        return new ArrayList<>(dedupByOrderId.values());
    }

    private String safeTrim(String s) {
        return s == null ? "" : s.trim();
    }
}

src/main/java/com/example/datapipeline/etl/JdbcLoader.java

package com.example.datapipeline.etl;

import com.example.datapipeline.model.Order;
import com.example.datapipeline.model.User;
import com.zaxxer.hikari.HikariDataSource;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.sql.Connection;
import java.sql.PreparedStatement;
import java.util.List;

public class JdbcLoader {
    private static final Logger log = LoggerFactory.getLogger(JdbcLoader.class);

    private final HikariDataSource ds;

    public JdbcLoader(HikariDataSource ds) {
        this.ds = ds;
    }

    public int upsertUsers(List<User> users) {
        if (users.isEmpty()) return 0;
        String sql = """
            INSERT INTO users (user_id, email, full_name, signup_date, updated_at)
            VALUES (?, ?, ?, ?, NOW())
            ON CONFLICT (email) DO UPDATE SET
              user_id = EXCLUDED.user_id,
              full_name = EXCLUDED.full_name,
              signup_date = EXCLUDED.signup_date,
              updated_at = NOW()
            """;
        int[] counts;
        try (Connection c = ds.getConnection();
             PreparedStatement ps = c.prepareStatement(sql)) {
            for (User u : users) {
                if (u.userId() == null) ps.setObject(1, null);
                else ps.setLong(1, u.userId());
                ps.setString(2, u.email());
                ps.setString(3, u.fullName());
                if (u.signupDate() == null) ps.setObject(4, null);
                else ps.setObject(4, u.signupDate());
                ps.addBatch();
            }
            counts = ps.executeBatch();
            c.commit();
        } catch (Exception e) {
            throw new RuntimeException("Upsert users failed", e);
        }
        int total = sum(counts);
        log.info("Upserted {} users", total);
        return total;
    }

    public int upsertOrders(List<Order> orders) {
        if (orders.isEmpty()) return 0;
        String sql = """
            INSERT INTO orders (order_id, user_email, amount, currency, created_at, status, updated_at)
            VALUES (?, ?, ?, ?, ?, ?, NOW())
            ON CONFLICT (order_id) DO UPDATE SET
              user_email = EXCLUDED.user_email,
              amount = EXCLUDED.amount,
              currency = EXCLUDED.currency,
              created_at = EXCLUDED.created_at,
              status = EXCLUDED.status,
              updated_at = NOW()
            """;
        int[] counts;
        try (Connection c = ds.getConnection();
             PreparedStatement ps = c.prepareStatement(sql)) {
            for (Order o : orders) {
                ps.setLong(1, o.orderId());
                ps.setString(2, o.userEmail());
                ps.setBigDecimal(3, o.amount());
                ps.setString(4, o.currency());
                ps.setObject(5, o.createdAt());
                ps.setString(6, o.status());
                ps.addBatch();
            }
            counts = ps.executeBatch();
            c.commit();
        } catch (Exception e) {
            throw new RuntimeException("Upsert orders failed", e);
        }
        int total = sum(counts);
        log.info("Upserted {} orders", total);
        return total;
    }

    private int sum(int[] arr) {
        int s = 0;
        for (int i : arr) {
            // PostgreSQL JDBC returns SUCCESS_NO_INFO (-2) sometimes; treat as 1
            s += (i > 0 ? i : 1);
        }
        return s;
    }
}

src/main/java/com/example/datapipeline/model/RawUserRecord.java

package com.example.datapipeline.model;

public record RawUserRecord(String userId, String email, String fullName, String signupDate) {}

src/main/java/com/example/datapipeline/model/RawOrderRecord.java

package com.example.datapipeline.model;

public record RawOrderRecord(String orderId, String userEmail, String amount, String currency, String createdAt, String status) {}

src/main/java/com/example/datapipeline/model/User.java

package com.example.datapipeline.model;

import java.time.LocalDate;

public record User(Long userId, String email, String fullName, LocalDate signupDate) {}

src/main/java/com/example/datapipeline/model/Order.java

package com.example.datapipeline.model;

import java.math.BigDecimal;
import java.time.OffsetDateTime;

public record Order(Long orderId, String userEmail, BigDecimal amount, String currency, OffsetDateTime createdAt, String status) {}

src/main/java/com/example/datapipeline/util/EmailValidator.java

package com.example.datapipeline.util;

import java.util.regex.Pattern;

public class EmailValidator {
    // 简化的邮箱校验(足够用于演示)
    private static final Pattern P = Pattern.compile("^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\\.[A-Za-z]{2,}$");

    public static boolean isValid(String email) {
        return email != null && P.matcher(email).matches();
    }
}

src/test/java/com/example/datapipeline/etl/TransformServiceTest.java

package com.example.datapipeline.etl;

import com.example.datapipeline.model.*;
import org.junit.jupiter.api.Test;

import java.util.List;

import static org.junit.jupiter.api.Assertions.*;

public class TransformServiceTest {

    @Test
    void testTransformUsers() {
        TransformService t = new TransformService();
        List<RawUserRecord> raw = List.of(
                new RawUserRecord("1", " Alice@Example.com ", "Alice", "2023-01-01"),
                new RawUserRecord("bad", "invalid", "Bad", "2023-02-02"),
                new RawUserRecord("2", "alice@example.com", "Alice Z", "2023-01-01")
        );
        var users = t.transformUsers(raw);
        assertEquals(1, users.size());
        assertEquals("alice@example.com", users.get(0).email());
        assertEquals("Alice Z", users.get(0).fullName());
    }

    @Test
    void testTransformOrders() {
        TransformService t = new TransformService();
        var users = t.transformUsers(List.of(
                new RawUserRecord("1", "alice@example.com", "Alice", "2023-01-01")
        ));
        var orders = t.transformOrders(List.of(
                new RawOrderRecord("1001", "alice@example.com", "10.0", "usd", "2023-12-31T23:59:59Z", "paid"),
                new RawOrderRecord("bad", "alice@example.com", "10", "USD", "2023-01-01T00:00:00Z", "PAID"),
                new RawOrderRecord("1001", "alice@example.com", "12.34", "USD", "2023-12-31T23:59:59Z", "REFUND")
        ), users);

        assertEquals(1, orders.size());
        assertEquals(1001L, orders.get(0).orderId());
        assertEquals("REFUND", orders.get(0).status());
    }
}

data/input/users.csv

user_id,email,full_name,signup_date
1, Alice@example.com ,Alice Zhang,2023-01-10
2,bob@example.com,Bob Li,2023-03-22
3,invalid-email,Bad User,2023-05-01
4,carol@example.com,  Carol Wang  ,2023-06-15
2,bob@example.com,Bob Li Jr,2023-03-22

data/input/orders.csv

order_id,user_email,amount,currency,created_at,status
1001,alice@example.com,19.99,USD,2023-12-31T23:59:59Z,PAID
1002, bob@example.com ,29.50,USD,2024-01-01T08:00:00+08:00,PAID
1003,carol@example.com,12.00,USD,2024-01-02T10:30:00Z,CANCELLED
1004,unknown@example.com,5.00,USD,2024-01-03T10:30:00Z,PAID
1002,bob@example.com,30.00,USD,2024-01-05T09:00:00Z,REFUND
badrow,bademail,xx,USD,not-a-time,PENDING

简历可写描述(技术栈与工作内容)

项目名称:基于 Java 的 CSV 至 PostgreSQL ETL 数据入门项目

  • 技术栈

    • 语言:Java 17
    • 构建:Maven,Shade Plugin(可执行 fat jar)
    • 数据库:PostgreSQL(Docker Compose 本地部署),Adminer(可视化)
    • 连接池:HikariCP
    • 数据迁移:Flyway
    • CSV 解析:OpenCSV
    • 日志:SLF4J + Logback
    • 测试:JUnit 5
    • 部署/运行:Docker Compose,一键启动数据库环境
  • 负责/实现

    • 设计并实现端到端 ETL 流程:从 CSV 抽取,基础清洗转换(字段标准化、邮箱校验、时间/金额解析、去重),批量写入 PostgreSQL;
    • 构建 JDBC 批处理 Upsert(ON CONFLICT)以保证幂等更新与数据一致性;
    • 使用 Flyway 管理数据库 schema 版本,保障环境可重复搭建;
    • 应用层可配置化(application.properties + 环境变量覆盖),支持多环境切换;
    • 引入连接池与批处理提升加载性能,并通过日志记录数据质检结果(有效/丢弃条数);
    • 编写单元测试验证转换逻辑正确性,使用 Docker Compose 提供本地开发数据库环境。
  • 产出与价值

    • 快速搭建入门级数据集成框架,可扩展到更多数据源/目标;
    • 具备幂等、可观测、可迁移特性,适合作为数据工程入门项目模板。

可扩展方向

  • 增量调度:引入 Quartz/Spring Scheduler,基于文件到达或时间窗增量同步
  • 更丰富数据源:REST API、JSON、Parquet、Kafka
  • 任务编排与监控:Airflow/Dagster + Prometheus/Grafana
  • 更严格的数据质量:Great Expectations/Deequ;错误数据落盘隔离
  • 更完善 DDL:分区表、物化视图、索引优化、权限分离(只读/写入用户)
  • CI/CD:GitHub Actions + Docker 镜像构建 + 环境部署
Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐