暂无图片
暂无图片
2
暂无图片
暂无图片
暂无图片

【金仓数据库征文】MyBatis-Plus+Spring Batch双剑合璧——金仓数据库批处理实战

原创 曾云林 2026-07-23
963

一、Spring Batch是真的香

我们单位有套老的报表系统,每天晚上跑批,处理几百万条数据,以前是跑在Oracle上的。信创改造要换成金仓,外包团队说Spring Batch连金仓有问题,搞不定,建议把批处理全部重写成存储过程。我一听就头大。那套批处理逻辑有好几十个job,上千行的配置,全改成存储过程?那得改到猴年马月去,而且出了问题排查更麻烦。
我心里清楚,Spring Batch本身是不挑数据库的,只要JDBC驱动没问题,理论上就能跑。问题出在方言适配、分页查询、序列生成这些细节上。MyBatis-Plus从v3.3.0开始正式宣布支持金仓了,这是个好消息,但Spring Batch呢?官方文档里没提过金仓的事。于是我就把这个当成了挑战,一边做项目一边研究,前前后后踩了几周的坑,终于把整套批处理系统从Oracle平滑迁到了金仓上。这篇文章就是全过程的实战记录,全是干货,代码拷过去就能用。Spring Batch批处理框架和金仓的适配之后,是真的香!

先交代环境:

  • 数据库:KingbaseES V9R3C18 MySQL兼容版
  • 端口:54321
  • 驱动:com.kingbase8.Driver
  • Java版本:JDK 17
  • Spring Boot:2.7.x
  • MyBatis-Plus:3.5.x(v3.3.0+官方支持金仓)
  • Spring Batch:4.3.x

二、Java驱动兼容适配实战

2.1 Maven依赖配置

金仓的JDBC驱动在官方Maven仓库里有,也可以本地安装。先把依赖配上:

image.png

<!-- pom.xml --> <dependencies> <!-- 金仓JDBC驱动 --> <dependency> <groupId>cn.com.kingbase</groupId> <artifactId>kingbase8</artifactId> <version>8.6.0</version> </dependency> <!-- MyBatis-Plus Boot Starter --> <dependency> <groupId>com.baomidou</groupId> <artifactId>mybatis-plus-boot-starter</artifactId> <version>3.5.3.1</version> </dependency> <!-- Spring Batch --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-batch</artifactId> </dependency> <!-- HikariCP连接池(Spring Boot默认带) --> <dependency> <groupId>com.zaxxer</groupId> <artifactId>HikariCP</artifactId> </dependency> </dependencies>

踩坑记录1:驱动包名的变迁

老版本的金仓驱动包名是com.kingbase.Driver,从V8开始改成了com.kingbase8.Driver,Maven的groupId也变过。如果你在网上搜到的教程里是com.kingbase.Driver,那大概率是老版本的,V9R3C18要用com.kingbase8.Driver

我最开始就踩了这个坑,配的老驱动类名,启动报ClassNotFoundException,查了半天。

2.2 连接测试

先写个最简单的连接测试,确认驱动和数据库都没问题:

image.png

package com.example.kingbase; import java.sql.Connection; import java.sql.DriverManager; import java.sql.ResultSet; import java.sql.Statement; public class KingbaseConnectionTest { public static void main(String[] args) { // MySQL兼容模式连接串 String url = "jdbc:kingbase8://127.0.0.1:54321/test?currentSchema=public&characterEncoding=utf8"; String username = "system"; String password = "123456"; try { // 加载驱动 Class.forName("com.kingbase8.Driver"); System.out.println("正在连接金仓数据库..."); Connection conn = DriverManager.getConnection(url, username, password); System.out.println("连接成功!"); // 测试查询 Statement stmt = conn.createStatement(); ResultSet rs = stmt.executeQuery("SELECT version()"); if (rs.next()) { System.out.println("版本信息:" + rs.getString(1)); } // 测试MySQL兼容模式 rs = stmt.executeQuery("SHOW SERVER_VERSION"); if (rs.next()) { System.out.println("SERVER_VERSION: " + rs.getString(1)); } rs.close(); stmt.close(); conn.close(); System.out.println("连接已关闭。"); } catch (Exception e) { System.err.println("连接失败:" + e.getMessage()); e.printStackTrace(); } } }

踩坑记录2:连接串参数

金仓的JDBC URL格式是jdbc:kingbase8://host:port/dbname,和PostgreSQL的jdbc:postgresql://格式类似,只是前缀不同。

常用参数:

  • currentSchema:指定默认schema
  • characterEncoding:字符集,建议utf8
  • useUnicode:是否使用Unicode
  • loginTimeout:登录超时时间(秒)

2.3 MySQL兼容模式关键注意事项

这个必须单独拿出来说,踩的坑太多了。

坑一:||运算符不是字符串拼接

这是最经典的坑,没有之一。MySQL兼容模式下,||是逻辑或运算符,不是字符串连接符。

-- 错误写法!返回布尔值,不是字符串拼接 SELECT username || ' - ' || real_name FROM sys_user; -- 正确写法 SELECT CONCAT(username, ' - ', real_name) FROM sys_user;

我第一次写报表SQL的时候,拼接字段用了||,结果返回的全是t,愣了半天没反应过来。后来查了金仓文档,MySQL兼容模式下sql_mode默认设成了和MySQL一致的行为,||就变成了逻辑或。

坑二:自增主键的获取方式

MySQL兼容模式下,LAST_INSERT_ID()函数是可用的,但用JDBC的getGeneratedKeys也能拿到。MyBatis-Plus已经帮我们处理好了,但如果你手写原生JDBC,要注意一下。

坑三:分页查询方言

金仓用的是PG风格的LIMIT ... OFFSET ...,MySQL兼容模式下也支持LIMIT m, n的写法。但Spring Batch的PagingQueryProvider需要正确的方言配置,不然分页SQL会生成错。


三、连接池调优实战(HikariCP)

Spring Boot默认用HikariCP,号称"最快的连接池"。默认配置其实就不错,但生产环境还是得根据实际情况调一调。

3.1 基础配置

image.png

# application.yml spring: datasource: driver-class-name: com.kingbase8.Driver url: jdbc:kingbase8://127.0.0.1:54321/test?currentSchema=public&characterEncoding=utf8 username: system password: 123456 # HikariCP 配置 hikari: # 连接池名称 pool-name: KingbaseHikariCP # 最小空闲连接数 minimum-idle: 5 # 最大连接池大小 maximum-pool-size: 20 # 连接超时时间(毫秒) connection-timeout: 30000 # 空闲连接超时时间(毫秒) idle-timeout: 600000 # 连接最大存活时间(毫秒) max-lifetime: 1800000 # 连接测试查询 connection-test-query: SELECT 1 # 连接初始化SQL connection-init-sql: SET NAMES utf8

3.2 连接池大小怎么算

这个问题很多人问,我统一说一下我的经验:

计算公式:连接数 = ((核心数 * 2) + 有效磁盘数)

这个公式是PostgreSQL官方给的,金仓同样适用。比如8核CPU、1块SSD盘,那就是 8*2 + 1 = 17,凑个整20就差不多了。

但这只是理论值,实际还要看你的业务类型:

  • CPU密集型:连接数不用太多,接近核心数就行
  • IO密集型:连接数可以多一些,因为很多连接在等IO
  • 混合场景:一般20-50个连接足够了

千万别觉得连接池越大越好。连接太多,数据库端的进程/线程上下文切换开销会急剧增大,反而拖累性能。我见过有人把连接池设成200的,结果数据库CPU飙得老高,一查全是空闲连接占着资源。

3.3 连接泄漏排查

HikariCP有个非常好用的功能——连接泄漏检测:

spring: datasource: hikari: # 连接占用超过这个时间就打日志告警(毫秒) leak-detection-threshold: 60000

设成60秒,意思是如果一个连接被借出去超过60秒还没还,就打一条告警日志。这个功能帮我抓了好几次连接泄漏——都是代码里开了事务忘了关,或者某个查询跑太久了。

3.4 生产环境监控

光配置好不行,还得能监控。HikariCP自带JMX监控,可以接Prometheus或者直接看日志。我一般配合Micrometer用:

image.png

package com.example.kingbase.config; import com.zaxxer.hikari.HikariDataSource; import io.micrometer.core.instrument.MeterRegistry; import org.springframework.context.annotation.Configuration; import javax.annotation.PostConstruct; import javax.sql.DataSource; @Configuration public class HikariMetricsConfig { private final DataSource dataSource; private final MeterRegistry meterRegistry; public HikariMetricsConfig(DataSource dataSource, MeterRegistry meterRegistry) { this.dataSource = dataSource; this.meterRegistry = meterRegistry; } @PostConstruct public void bindMetrics() { if (dataSource instanceof HikariDataSource) { HikariDataSource hikari = (HikariDataSource) dataSource; hikari.setMetricRegistry(meterRegistry); } } }

然后Grafana看板上就能看到这些指标:

  • 活跃连接数
  • 空闲连接数
  • 等待连接的线程数
  • 连接获取时间
  • 连接创建总数

有了这些数据,调优就不是拍脑袋了。

3.5 常见问题排查

问题:连接获取超时

  • 现象:HikariPool-1 - Connection is not available, request timed out after 30000ms
  • 排查思路:
    1. 是不是连接泄漏了?开leak-detection-threshold看看
    2. 是不是SQL太慢,连接占着不放?
    3. 连接池是不是太小了?看等待连接的线程数
    4. 数据库端是不是有锁等待?查sys_locks或者活动会话

问题:连接频繁断开

  • 现象:报连接已关闭、I/O错误之类的
  • 排查思路:
    1. 是不是防火墙把空闲连接掐了?
    2. 数据库的idle_in_transaction_session_timeout是不是设太短?
    3. max-lifetime要比数据库的连接超时时间短
    4. 可以适当调小idle-timeout,让空闲连接及时释放

四、MyBatis-Plus集成金仓实战

好消息:MyBatis-Plus从v3.3.0开始正式支持金仓了!不用自己写方言,不用自己适配,开箱即用。

4.1 配置

image.png

# application.yml mybatis-plus: mapper-locations: classpath*:mapper/**/*.xml type-aliases-package: com.example.kingbase.entity configuration: # 驼峰命名转换 map-underscore-to-camel-case: true # 日志 log-impl: org.apache.ibatis.logging.slf4j.Slf4jImpl global-config: db-config: # 主键类型:自增(数据库BIGSERIAL) id-type: auto # 逻辑删除 logic-delete-field: deleted logic-delete-value: 1 logic-not-delete-value: 0

踩坑记录3:主键策略

金仓有好几种主键生成方式:

  • BIGSERIAL(序列,推荐)
  • UUID
  • 雪花算法(应用层生成)

如果你用BIGSERIAL,那id-type设成auto就行,MyBatis-Plus会自动处理,插入后主键会回填到实体对象里。

如果你用的是MySQL兼容模式的AUTO_INCREMENT,也能正常工作,本质上金仓内部还是用序列实现的。

4.2 实体类定义

image.png

package com.example.kingbase.entity; import com.baomidou.mybatisplus.annotation.*; import lombok.Data; import java.time.LocalDateTime; @Data @TableName("sys_user") public class SysUser { @TableId(type = IdType.AUTO) private Long id; private String username; private String password; private String realName; private String email; private String phone; /** * 状态:1-启用,0-禁用 */ private Integer status; /** * 创建时间 */ @TableField(fill = FieldFill.INSERT) private LocalDateTime createTime; /** * 更新时间 */ @TableField(fill = FieldFill.INSERT_UPDATE) private LocalDateTime updateTime; /** * 逻辑删除标记 */ @TableLogic private Integer deleted; }

4.3 Mapper和Service

image.png

package com.example.kingbase.mapper; import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.example.kingbase.entity.SysUser; import org.apache.ibatis.annotations.Mapper; @Mapper public interface SysUserMapper extends BaseMapper<SysUser> { // 啥都不用写,基础CRUD全有了 }

image.png

package com.example.kingbase.service.impl; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.core.metadata.IPage; import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import com.example.kingbase.entity.SysUser; import com.example.kingbase.mapper.SysUserMapper; import com.example.kingbase.service.SysUserService; import org.springframework.stereotype.Service; @Service public class SysUserServiceImpl extends ServiceImpl<SysUserMapper, SysUser> implements SysUserService { /** * 分页查询启用的用户 */ @Override public IPage<SysUser> pageActiveUsers(int pageNum, int pageSize, String keyword) { Page<SysUser> page = new Page<>(pageNum, pageSize); LambdaQueryWrapper<SysUser> wrapper = new LambdaQueryWrapper<>(); wrapper.eq(SysUser::getStatus, 1) .like(keyword != null && !keyword.isEmpty(), SysUser::getUsername, keyword) .orderByDesc(SysUser::getCreateTime); return this.page(page, wrapper); } }

4.4 分页插件配置

MyBatis-Plus的分页插件需要配置方言。好消息是,v3.3.0+已经内置了金仓的方言:

image.png

package com.example.kingbase.config; import com.baomidou.mybatisplus.annotation.DbType; import com.baomidou.mybatisplus.extension.plugins.MybatisPlusInterceptor; import com.baomidou.mybatisplus.extension.plugins.inner.PaginationInnerInterceptor; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class MyBatisPlusConfig { @Bean public MybatisPlusInterceptor mybatisPlusInterceptor() { MybatisPlusInterceptor interceptor = new MybatisPlusInterceptor(); // 分页插件,指定数据库类型为金仓 interceptor.addInnerInterceptor(new PaginationInnerInterceptor(DbType.KINGBASE_ES)); return interceptor; } }

划重点: DbType.KINGBASE_ES,这个枚举值就是v3.3.0加进来的。以前版本没有的话,你可能得自己写方言类,现在不用了,官方支持。

踩坑记录4:分页总数查询性能

MyBatis-Plus分页会自动生成一条SELECT COUNT(1) FROM (...)的count查询。如果你的表很大、查询条件又复杂,这条count查询可能很慢。

优化思路:

  1. 确保WHERE条件里的字段有索引
  2. 如果不需要精确总数(比如前端只需要"加载更多"),可以关掉count查询:
Page<SysUser> page = new Page<>(pageNum, pageSize, false); // 第三个参数false=不查总数
  1. 超大数据量表考虑用游标查询(Spring Batch里常用)

4.5 批量操作优化

MyBatis-Plus的saveBatch默认是一条条INSERT的,效率不行。大数据量插入要开批量模式:

image.png

// 方式一:MyBatis-Plus自带的批量插入(需要配置) // 方式二:用SqlSession的BATCH模式 @Autowired private SqlSessionFactory sqlSessionFactory; public void batchInsertUsers(List<SysUser> users) { // 开启BATCH模式的SqlSession try (SqlSession sqlSession = sqlSessionFactory.openSession(ExecutorType.BATCH)) { SysUserMapper mapper = sqlSession.getMapper(SysUserMapper.class); for (int i = 0; i < users.size(); i++) { mapper.insert(users.get(i)); // 每1000条提交一次,避免内存溢出 if (i > 0 && i % 1000 == 0) { sqlSession.commit(); } } sqlSession.commit(); } }

但说实话,BATCH模式虽然比一条条快,但和真正的批量INSERT还是有差距。如果数据量特别大(十万级以上),我更推荐用金仓的COPY命令,或者Spring Batch的Chunk模式+JDBC批处理。


五、Spring Batch集成金仓实战

这是这篇文章的重头戏。只要搞清楚原理,适配起来并不难。

5.1 Spring Batch的数据库依赖

Spring Batch本身需要一套元数据表(BATCH_JOB_INSTANCE、BATCH_JOB_EXECUTION等),用来记录Job的执行状态、参数、步骤信息。这些表的创建脚本在Spring Batch的jar包里,每种数据库有对应的脚本。两种思路:

方案一:用PostgreSQL的脚本(推荐)

金仓基于PostgreSQL内核,语法高度兼容。直接拿PostgreSQL的脚本来用,基本没问题。

方案二:用内存数据库存元数据

Spring Batch元数据存H2之类的内存数据库,业务数据存在金仓。适合非生产环境或者简单场景。

我推荐方案一,生产环境更可靠。

5.2 元数据表初始化

先把PostgreSQL的脚本找出来,改一改就能用了。脚本位置在spring-batch-core.jar里的org/springframework/batch/core/schema-postgresql.sql

我直接把改好的金仓版贴出来:

-- 金仓版 Spring Batch 元数据建表脚本 -- 基于PostgreSQL脚本修改,MySQL兼容模式下可正常运行 CREATE TABLE BATCH_JOB_INSTANCE ( JOB_INSTANCE_ID BIGINT NOT NULL PRIMARY KEY, VERSION BIGINT, JOB_NAME VARCHAR(100) NOT NULL, JOB_KEY VARCHAR(32) NOT NULL, CONSTRAINT JOB_INST_UN UNIQUE (JOB_NAME, JOB_KEY) ); CREATE TABLE BATCH_JOB_EXECUTION ( JOB_EXECUTION_ID BIGINT NOT NULL PRIMARY KEY, VERSION BIGINT, JOB_INSTANCE_ID BIGINT NOT NULL, CREATE_TIME TIMESTAMP NOT NULL, START_TIME TIMESTAMP DEFAULT NULL, END_TIME TIMESTAMP DEFAULT NULL, STATUS VARCHAR(10), EXIT_CODE VARCHAR(2500), EXIT_MESSAGE VARCHAR(2500), LAST_UPDATED TIMESTAMP, JOB_CONFIGURATION_LOCATION VARCHAR(2500) NULL ); CREATE TABLE BATCH_JOB_EXECUTION_PARAMS ( JOB_EXECUTION_ID BIGINT NOT NULL, TYPE_CD VARCHAR(6) NOT NULL, KEY_NAME VARCHAR(100) NOT NULL, STRING_VAL VARCHAR(250), DATE_VAL TIMESTAMP DEFAULT NULL, LONG_VAL BIGINT, DOUBLE_VAL DOUBLE PRECISION, IDENTIFYING CHAR(1) NOT NULL, CONSTRAINT JOB_EXEC_PARAMS_FK FOREIGN KEY (JOB_EXECUTION_ID) REFERENCES BATCH_JOB_EXECUTION(JOB_EXECUTION_ID) ); CREATE TABLE BATCH_STEP_EXECUTION ( STEP_EXECUTION_ID BIGINT NOT NULL PRIMARY KEY, VERSION BIGINT NOT NULL, STEP_NAME VARCHAR(100) NOT NULL, JOB_EXECUTION_ID BIGINT NOT NULL, START_TIME TIMESTAMP NOT NULL, END_TIME TIMESTAMP DEFAULT NULL, STATUS VARCHAR(10), COMMIT_COUNT BIGINT, READ_COUNT BIGINT, FILTER_COUNT BIGINT, WRITE_COUNT BIGINT, READ_SKIP_COUNT BIGINT, WRITE_SKIP_COUNT BIGINT, PROCESS_SKIP_COUNT BIGINT, ROLLBACK_COUNT BIGINT, EXIT_CODE VARCHAR(2500), EXIT_MESSAGE VARCHAR(2500), LAST_UPDATED TIMESTAMP, CONSTRAINT STEP_EXECUTION_FK FOREIGN KEY (JOB_EXECUTION_ID) REFERENCES BATCH_JOB_EXECUTION(JOB_EXECUTION_ID) ); CREATE TABLE BATCH_STEP_EXECUTION_CONTEXT ( STEP_EXECUTION_ID BIGINT NOT NULL PRIMARY KEY, SHORT_CONTEXT VARCHAR(2500) NOT NULL, SERIALIZED_CONTEXT TEXT, CONSTRAINT STEP_EXEC_CTX_FK FOREIGN KEY (STEP_EXECUTION_ID) REFERENCES BATCH_STEP_EXECUTION(STEP_EXECUTION_ID) ); CREATE TABLE BATCH_JOB_EXECUTION_CONTEXT ( JOB_EXECUTION_ID BIGINT NOT NULL PRIMARY KEY, SHORT_CONTEXT VARCHAR(2500) NOT NULL, SERIALIZED_CONTEXT TEXT, CONSTRAINT JOB_EXEC_CTX_FK FOREIGN KEY (JOB_EXECUTION_ID) REFERENCES BATCH_JOB_EXECUTION(JOB_EXECUTION_ID) ); -- 序列 CREATE SEQUENCE BATCH_JOB_SEQ START WITH 1 INCREMENT BY 1 NO MINVALUE NO MAXVALUE CACHE 20; CREATE SEQUENCE BATCH_JOB_EXECUTION_SEQ START WITH 1 INCREMENT BY 1 NO MINVALUE NO MAXVALUE CACHE 20; CREATE SEQUENCE BATCH_STEP_EXECUTION_SEQ START WITH 1 INCREMENT BY 1 NO MINVALUE NO MAXVALUE CACHE 20;

金仓支持SEQUENCE,和PostgreSQL语法一样,所以序列这块直接用就行。

踩坑记录5:序列名的大小写

有个细节要注意,Spring Batch默认查序列的时候用的是大写的序列名。金仓默认是大小写不敏感的(会转成小写),如果序列创建的时候是小写的,查询可能找不到。

保险起见,建序列的时候就用大写名字,或者设置Spring Batch的table-prefix和数据库里的一致。

5.3 Spring Batch配置类

image.png

package com.example.kingbase.batch.config; import org.springframework.batch.core.configuration.annotation.DefaultBatchConfigurer; import org.springframework.batch.core.configuration.annotation.EnableBatchProcessing; import org.springframework.batch.core.launch.JobLauncher; import org.springframework.batch.core.launch.support.SimpleJobLauncher; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.repository.support.JobRepositoryFactoryBean; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.transaction.PlatformTransactionManager; import javax.sql.DataSource; @Configuration @EnableBatchProcessing public class BatchConfig extends DefaultBatchConfigurer { private final DataSource dataSource; private final PlatformTransactionManager transactionManager; public BatchConfig(DataSource dataSource, PlatformTransactionManager transactionManager) { this.dataSource = dataSource; this.transactionManager = transactionManager; } @Override protected JobRepository createJobRepository() throws Exception { JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean(); factory.setDataSource(dataSource); factory.setTransactionManager(transactionManager); // 关键:指定数据库类型为POSTGRES // 金仓基于PG内核,用PG的SQL类型完全兼容 factory.setDatabaseType("POSTGRES"); // 设置表前缀(和你建表时的前缀一致) factory.setTablePrefix("BATCH_"); factory.setIsolationLevelForCreate("ISOLATION_READ_COMMITTED"); factory.setMaxVarCharLength(2500); return factory.getObject(); } @Bean @Override public JobLauncher createJobLauncher() throws Exception { SimpleJobLauncher jobLauncher = new SimpleJobLauncher(); jobLauncher.setJobRepository(createJobRepository()); jobLauncher.afterPropertiesSet(); return jobLauncher; } }

划重点: factory.setDatabaseType("POSTGRES"),这一行是关键。Spring Batch支持的数据库类型里没有金仓,但金仓和PostgreSQL高度兼容,所以直接指定为POSTGRES就能用。

我最开始没设这个,Spring Boot自动检测的时候识别成了别的类型,结果生成的SQL语法不对,启动报错。加了这一行就全好了。

5.4 实战:用户数据同步批处理Job

来一个真实场景的例子——从源表(比如MongoDB或者另一个库)同步用户数据到金仓的sys_user表。这是我项目里真实的需求,简化一下放出来。

image.png

package com.example.kingbase.batch.job; import com.example.kingbase.entity.SysUser; import com.example.kingbatch.mapper.SysUserMapper; import org.springframework.batch.core.Job; import org.springframework.batch.core.Step; import org.springframework.batch.core.configuration.annotation.JobBuilderFactory; import org.springframework.batch.core.configuration.annotation.StepBuilderFactory; import org.springframework.batch.core.launch.support.RunIdIncrementer; import org.springframework.batch.item.ItemProcessor; import org.springframework.batch.item.ItemReader; import org.springframework.batch.item.ItemWriter; import org.springframework.batch.item.database.JdbcPagingItemReader; import org.springframework.batch.item.database.Order; import org.springframework.batch.item.database.support.PostgresPagingQueryProvider; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import javax.sql.DataSource; import java.sql.ResultSet; import java.util.HashMap; import java.util.Map; @Configuration public class UserSyncJobConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final DataSource dataSource; private final SysUserMapper sysUserMapper; public UserSyncJobConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, DataSource dataSource, SysUserMapper sysUserMapper) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.dataSource = dataSource; this.sysUserMapper = sysUserMapper; } /** * 用户数据同步Job * 场景:从用户原始数据表同步到处理后的用户表 */ @Bean public Job userSyncJob() { return jobBuilderFactory.get("userSyncJob") .incrementer(new RunIdIncrementer()) .start(userSyncStep()) .build(); } @Bean public Step userSyncStep() { return stepBuilderFactory.get("userSyncStep") .<SysUser, SysUser>chunk(100) // 每100条提交一次 .reader(userSyncReader()) .processor(userSyncProcessor()) .writer(userSyncWriter()) .faultTolerant() .skip(Exception.class) .skipLimit(10) // 最多跳过10条异常数据 .listener(new UserSyncStepListener()) .build(); } /** * 读取器:分页读取源数据 * 用PostgresPagingQueryProvider,金仓完全兼容 */ @Bean public ItemReader<SysUser> userSyncReader() { JdbcPagingItemReader<SysUser> reader = new JdbcPagingItemReader<>(); reader.setDataSource(dataSource); reader.setPageSize(100); reader.setFetchSize(100); // 关键:用PostgreSQL的分页查询提供者 PostgresPagingQueryProvider queryProvider = new PostgresPagingQueryProvider(); queryProvider.setSelectClause("SELECT id, username, real_name, email, phone, status, create_time"); queryProvider.setFromClause("FROM user_source"); queryProvider.setWhereClause("WHERE sync_status = 0"); // 排序键(分页必须有排序) Map<String, Order> sortKeys = new HashMap<>(); sortKeys.put("id", Order.ASCENDING); queryProvider.setSortKeys(sortKeys); reader.setQueryProvider(queryProvider); // 行映射 reader.setRowMapper((ResultSet rs, int rowNum) -> { SysUser user = new SysUser(); user.setId(rs.getLong("id")); user.setUsername(rs.getString("username")); user.setRealName(rs.getString("real_name")); user.setEmail(rs.getString("email")); user.setPhone(rs.getString("phone")); user.setStatus(rs.getInt("status")); return user; }); return reader; } /** * 处理器:数据清洗、转换 */ @Bean public ItemProcessor<SysUser, SysUser> userSyncProcessor() { return user -> { // 示例:手机号脱敏、密码加密、邮箱转小写等 if (user.getEmail() != null) { user.setEmail(user.getEmail().toLowerCase()); } // 设置默认密码(真实场景用加密) user.setPassword("e10adc3949ba59abbe56e057f20f883e"); return user; }; } /** * 写入器:批量写入金仓 * 用MyBatis-Plus的批量写入 */ @Bean public ItemWriter<SysUser> userSyncWriter() { return items -> { // chunk模式下,items就是这一批的数据 // 批量插入或更新 for (SysUser user : items) { // 存在则更新,不存在则插入 sysUserMapper.insertOrUpdate(user); } }; } }

5.5 Step执行监听器

image.png

package com.example.kingbase.batch.job; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.batch.core.ExitStatus; import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.StepExecutionListener; public class UserSyncStepListener implements StepExecutionListener { private static final Logger log = LoggerFactory.getLogger(UserSyncStepListener.class); private long startTime; @Override public void beforeStep(StepExecution stepExecution) { startTime = System.currentTimeMillis(); log.info("【用户同步步骤】开始执行,Step: {}", stepExecution.getStepName()); } @Override public ExitStatus afterStep(StepExecution stepExecution) { long duration = System.currentTimeMillis() - startTime; log.info("【用户同步步骤】执行完成,耗时: {}ms", duration); log.info(" 读取数: {}", stepExecution.getReadCount()); log.info(" 写入数: {}", stepExecution.getWriteCount()); log.info(" 跳过数: {}", stepExecution.getReadSkipCount() + stepExecution.getWriteSkipCount()); log.info(" 提交次数: {}", stepExecution.getCommitCount()); log.info(" 执行状态: {}", stepExecution.getStatus()); return stepExecution.getExitStatus(); } }

5.6 触发Job执行

image.png

package com.example.kingbase.batch.controller; import org.springframework.batch.core.Job; import org.springframework.batch.core.JobParameters; import org.springframework.batch.core.JobParametersBuilder; import org.springframework.batch.core.launch.JobLauncher; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; @RestController @RequestMapping("/batch") public class BatchController { @Resource private JobLauncher jobLauncher; @Resource(name = "userSyncJob") private Job userSyncJob; @PostMapping("/user-sync") public String runUserSyncJob() { try { JobParameters params = new JobParametersBuilder() .addLong("runTime", System.currentTimeMillis()) .toJobParameters(); jobLauncher.run(userSyncJob, params); return "用户同步Job已启动"; } catch (Exception e) { return "启动失败: " + e.getMessage(); } } }

5.7 Spring Batch性能优化经验

优化一:Chunk大小调优

Chunk大小决定了每批处理多少条数据提交一次事务。太小了频繁提交影响性能,太大了占内存、回滚代价高。

我的经验值:

  • 简单的单表插入/更新:chunk=500~1000
  • 复杂的多表操作、有外部调用:chunk=50~100
  • 纯数据搬运(读-写):chunk=1000~5000

金仓的写入性能相当不错,我测试过单表批量插入,chunk=1000的时候,每秒能插3000+条。

优化二:JDBC批量写入

如果ItemWriter里用的是JDBC,可以开批量更新:

// 配置MyBatis的批量执行器 // 或者用JdbcBatchItemWriter @Bean public JdbcBatchItemWriter<SysUser> jdbcBatchWriter(DataSource dataSource) { JdbcBatchItemWriter<SysUser> writer = new JdbcBatchItemWriter<>(); writer.setDataSource(dataSource); writer.setSql("INSERT INTO sys_user (username, password, real_name, email, status) " + "VALUES (?, ?, ?, ?, ?)"); writer.setItemPreparedStatementSetter((item, ps) -> { ps.setString(1, item.getUsername()); ps.setString(2, item.getPassword()); ps.setString(3, item.getRealName()); ps.setString(4, item.getEmail()); ps.setInt(5, item.getStatus()); }); return writer; }

JdbcBatchItemWriter底层用的是JDBC的addBatch()executeBatch(),比循环insert快不少。

优化三:游标读取 vs 分页读取

数据量特别大的时候,JdbcPagingItemReader每次查一页,反复翻页性能会下降。这时候可以用游标读取:

@Bean public JdbcCursorItemReader<SysUser> cursorReader(DataSource dataSource) { JdbcCursorItemReader<SysUser> reader = new JdbcCursorItemReader<>(); reader.setDataSource(dataSource); reader.setSql("SELECT id, username, real_name FROM user_source ORDER BY id"); reader.setRowMapper(new BeanPropertyRowMapper<>(SysUser.class)); reader.setFetchSize(1000); // 游标获取批次大小 return reader; }

游标读取的好处是一次性打开游标,流式拉取数据,不用反复执行分页SQL。但要注意,游标打开期间数据库连接会一直占用,连接池要配够。

优化四:多线程并行处理

Spring Batch支持多线程Step,用TaskExecutor实现:

@Bean public Step userSyncStep() { return stepBuilderFactory.get("userSyncStep") .<SysUser, SysUser>chunk(100) .reader(userSyncReader()) .processor(userSyncProcessor()) .writer(userSyncWriter()) .taskExecutor(batchTaskExecutor()) .throttleLimit(5) // 并发线程数 .build(); } @Bean public TaskExecutor batchTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(100); executor.setThreadNamePrefix("batch-"); executor.initialize(); return executor; }

但注意,多线程模式下,ItemReader必须是线程安全的。分页读取JdbcPagingItemReader默认是线程安全的,但游标读取就不是了,需要加同步或者用分区Step。

5.8 金仓适配Spring Batch的坑点总结

坑点 现象 解决方案
元数据表创建失败 启动报BATCH_JOB_INSTANCE表不存在 手动执行PG版建表脚本
数据库类型不识别 JobRepository初始化报错 手动设置databaseType为"POSTGRES"
分页SQL语法错误 分页查询报语法错 用PostgresPagingQueryProvider
序列获取失败 生成ID时报序列不存在 序列名用大写,与Spring Batch默认一致
大字段存储异常 ExecutionContext存储失败 TEXT类型足够,不需要BLOB
字符串长度超限 VARCHAR长度不够报错 调整maxVarCharLength,金仓VARCHAR支持很大

六、开发大赛心得

6.1心得

纯JDBC连接金仓,写起来快,也不容易出错。但后来想想,那有什么意思呢?做了十几年DBA,这些不都是基本功吗。正好手里有那个报表系统信创改造的项目,Spring Batch连金仓是个实打实的硬骨头。外包团队啃不下来,我就想,我自己能不能啃下来?啃下来了,项目推进了,大赛成绩也有了,一举两得。

6.2 技术收获

这段时间做下来,收获还是挺大的。

第一,对MyBatis-Plus和Spring Batch的理解更深了。 以前就是用用,出了问题查资料。这次是从底层适配的角度去研究,方言怎么写的、分页怎么实现的、Spring Batch元数据表是干啥的、JobRepository是怎么存状态的。搞懂了原理,再遇到问题就不慌了。

第二,金仓的兼容性比我想象的好。 说实话,最开始我心里也打鼓,Spring Batch官方都没支持金仓,会不会到处是坑?真做下来发现,金仓和PostgreSQL的兼容度非常高,只要把databaseType设成POSTGRES,基本上都能跑通。核心的语法、函数、序列、事务,全都没问题。

第三,批处理性能优化的门道。 以前调SQL是一行一行调,现在是一批一批调。Chunk大小、fetch size、批量写入、多线程并行,每一项都对性能有影响。我最开始跑一个100万数据的同步job,要跑半个小时,优化完之后5分钟搞定,提升还是很明显的。

6.3 踩坑感悟

做开发类的题目,最大的感悟就是——隔行如隔山。

以前我总觉得,不就是写个Java程序吗,有什么难的。真自己上手写了才发现,框架那么多、配置那么杂、每个版本还不一样。Spring Batch的配置我前前后后调了好几天才跑通,一会儿是JobRepository没配置对,一会儿是事务管理器冲突,一会儿是分页方言不支持。

但反过来想,开发人员做DBA的活,估计也一样头大。索引怎么建、执行计划怎么看、锁等待怎么查,每一项都是学问。所以说,术业有专攻,DBA和开发团队配合好,才能把项目做好。

还有一个感悟就是——信创改造不是简单的替换数据库。表面上看就是改个连接串、换个驱动,但真正落地的时候,每一层框架、每一行SQL都可能出问题。有官方支持的(比如MyBatis-Plus)就好很多,没有官方支持的(比如Spring Batch)就得自己摸索。

6.4 对金仓生态的建议

最后说几点建议,都是真心话:

  1. 官方出一个Spring Boot Starter全家桶。 把金仓的数据源自动配置、方言、连接池默认参数都打包进去,开发者引入一个starter就完事了,不用自己去配这配那。
  2. Spring Batch官方适配可以安排一下。 现在用PostgreSQL模式能跑,但毕竟不是官方支持,心里总是不踏实。如果金仓能给Spring Batch贡献个官方的数据库类型,那就完美了。
  3. 批量导入工具再丰富些。 现在COPY命令很好用,但如果有类似MySQL的LOAD DATA LOCAL INFILE那种便捷的批量导入接口就更好了。
  4. 文档要跟上。 现在网上关于金仓的技术文章还是少,出了问题搜不到答案,只能自己摸。官方如果能出一套完整的开发指南,覆盖常用的框架和场景,开发者上手会快很多。

七、总结

这段时间,从最基础的JDBC连接,到MyBatis-Plus ORM,再到Spring Batch批处理框架,一层层往上走,踩了很多坑,也学到了很多东西。

回过头来看,金仓数据库的Java生态其实已经相当不错了:

  • 驱动层:金仓官方JDBC驱动成熟稳定,和PG高度兼容
  • ORM层:MyBatis-Plus从v3.3.0开始官方支持,开箱即用
  • 批处理层:Spring Batch用PostgreSQL模式可以完美运行

对于正在做信创改造的团队来说,Java技术栈迁移到金仓,技术上是完全可行的。关键是要把细节处理好——方言配置、SQL兼容性、连接池调优、批处理参数,每一项都要到位。

文章转载自曾云林,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论