如何用Spring R2DBC实现批量新增与更新未知状态的用户数据?
用Spring R2DBC实现批量用户数据的UPSERT操作
核心思路
针对百万级用户数据的插入/更新需求,优先采用数据库原生的UPSERT语法(不同数据库语法略有差异),避免先查询再判断的低效逻辑。结合Spring R2DBC的响应式批量操作能力,既能保证性能,又能适配大流量数据处理场景。
具体实现步骤
1. 定义用户实体类
先创建与数据库表映射的用户实体:
import org.springframework.data.annotation.Id; import org.springframework.data.relational.core.mapping.Table; @Table("users") public record User( @Id Long userId, String username, String email, Long lastUpdatedTime // 其他业务字段 ) {}
2. 基于数据库原生UPSERT实现批量操作
不同数据库的UPSERT语法不同,以下是主流数据库的实现示例:
PostgreSQL(ON CONFLICT语法)
通过DatabaseClient执行批量UPSERT:
import org.springframework.r2dbc.core.DatabaseClient; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Service public class UserBatchService { private final DatabaseClient databaseClient; public UserBatchService(DatabaseClient databaseClient) { this.databaseClient = databaseClient; } public Mono<Void> batchUpsertUsers(Flux<User> userFlux) { String sql = """ INSERT INTO users (user_id, username, email, last_updated_time) VALUES (:userId, :username, :email, :lastUpdatedTime) ON CONFLICT (user_id) DO UPDATE SET username = EXCLUDED.username, email = EXCLUDED.email, last_updated_time = EXCLUDED.last_updated_time """; return databaseClient.sql(sql) .filter((statement, executeFunction) -> statement.fetchSize(1000)) // 调整批量大小适配数据库性能 .bind("userId", User::userId) .bind("username", User::username) .bind("email", User::email) .bind("lastUpdatedTime", User::lastUpdatedTime) .thenMany(userFlux) .then(); } }
MySQL(ON DUPLICATE KEY UPDATE语法)
只需替换SQL语句即可:
String sql = """ INSERT INTO users (user_id, username, email, last_updated_time) VALUES (:userId, :username, :email, :lastUpdatedTime) ON DUPLICATE KEY UPDATE username = VALUES(username), email = VALUES(email), last_updated_time = VALUES(last_updated_time) """;
3. 百万级数据处理的关键注意事项
- 流式加载数据:读取文件时采用流式方式(比如
Files.lines()转Flux),避免一次性将100万条数据加载到内存导致溢出。 - 批量大小调优:根据数据库性能调整
fetchSize(通常设置1000-5000),平衡内存占用与数据库交互次数。 - 事务控制:如果需要保证数据一致性,可通过
TransactionalOperator添加事务支持:
import org.springframework.transaction.reactive.TransactionalOperator; public Mono<Void> batchUpsertUsersWithTransaction(Flux<User> userFlux) { return transactionalOperator.transactional( batchUpsertUsers(userFlux) ); }
4. 备选方案:先查后更(仅适合小批量数据)
如果数据库不支持UPSERT(极端场景),可采用先查询再判断的逻辑,但这种方式对百万级数据效率极低,不推荐:
public Mono<Void> upsertSingleUser(User user) { return databaseClient.select() .from(User.class) .matching(Criteria.where("userId").is(user.userId())) .one() .flatMap(existing -> databaseClient.update() .table(User.class) .using(user) .matching(Criteria.where("userId").is(user.userId())) .then()) .switchIfEmpty(databaseClient.insert() .into(User.class) .using(user) .then()); } // 批量处理时需控制并发数,避免数据库连接耗尽 public Mono<Void> batchUpsertUsersLowEfficiency(Flux<User> userFlux) { return userFlux.flatMap(this::upsertSingleUser, 10) .then(); }
总结
针对百万级用户数据,数据库原生UPSERT+Spring R2DBC批量操作是性能最优的方案,既避免了冗余查询,又能高效处理流式数据。需根据数据库类型调整UPSERT语法,并做好内存与并发控制。
内容的提问来源于stack exchange,提问作者Andrew Sukonin
相关产品推荐
相关产品推荐

