Java 数据开发入门项目:CSV -> Postgres 的小型ETL管道
一个可在本地运行的入门级数据开发项目:从 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 (本文即为说明)
快速开始
- 启动数据库(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
- 准备数据文件
- 已提供样例:data/input/users.csv、data/input/orders.csv
- 运行
- JDK 17+
- Maven 3.8+
- 先编译并打包可执行 jar:
mvn -q -DskipTests package
- 运行:
java -jar target/data-pipeline-basic-1.0.0-shaded.jar
- 查看结果
- 通过 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 镜像构建 + 环境部署
更多推荐


所有评论(0)