一、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仓库里有,也可以本地安装。先把依赖配上:
<!-- 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 连接测试
先写个最简单的连接测试,确认驱动和数据库都没问题:
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:指定默认schemacharacterEncoding:字符集,建议utf8useUnicode:是否使用UnicodeloginTimeout:登录超时时间(秒)
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 基础配置
# 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用:
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 - 排查思路:
- 是不是连接泄漏了?开leak-detection-threshold看看
- 是不是SQL太慢,连接占着不放?
- 连接池是不是太小了?看等待连接的线程数
- 数据库端是不是有锁等待?查
sys_locks或者活动会话
问题:连接频繁断开
- 现象:报连接已关闭、I/O错误之类的
- 排查思路:
- 是不是防火墙把空闲连接掐了?
- 数据库的
idle_in_transaction_session_timeout是不是设太短? max-lifetime要比数据库的连接超时时间短- 可以适当调小
idle-timeout,让空闲连接及时释放
四、MyBatis-Plus集成金仓实战
好消息:MyBatis-Plus从v3.3.0开始正式支持金仓了!不用自己写方言,不用自己适配,开箱即用。
4.1 配置
# 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 实体类定义
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
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全有了
}
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+已经内置了金仓的方言:
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查询可能很慢。
优化思路:
- 确保WHERE条件里的字段有索引
- 如果不需要精确总数(比如前端只需要"加载更多"),可以关掉count查询:
Page<SysUser> page = new Page<>(pageNum, pageSize, false); // 第三个参数false=不查总数
- 超大数据量表考虑用游标查询(Spring Batch里常用)
4.5 批量操作优化
MyBatis-Plus的saveBatch默认是一条条INSERT的,效率不行。大数据量插入要开批量模式:
// 方式一: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配置类
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表。这是我项目里真实的需求,简化一下放出来。
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执行监听器
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执行
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 对金仓生态的建议
最后说几点建议,都是真心话:
- 官方出一个Spring Boot Starter全家桶。 把金仓的数据源自动配置、方言、连接池默认参数都打包进去,开发者引入一个starter就完事了,不用自己去配这配那。
- Spring Batch官方适配可以安排一下。 现在用PostgreSQL模式能跑,但毕竟不是官方支持,心里总是不踏实。如果金仓能给Spring Batch贡献个官方的数据库类型,那就完美了。
- 批量导入工具再丰富些。 现在COPY命令很好用,但如果有类似MySQL的
LOAD DATA LOCAL INFILE那种便捷的批量导入接口就更好了。 - 文档要跟上。 现在网上关于金仓的技术文章还是少,出了问题搜不到答案,只能自己摸。官方如果能出一套完整的开发指南,覆盖常用的框架和场景,开发者上手会快很多。
七、总结
这段时间,从最基础的JDBC连接,到MyBatis-Plus ORM,再到Spring Batch批处理框架,一层层往上走,踩了很多坑,也学到了很多东西。
回过头来看,金仓数据库的Java生态其实已经相当不错了:
- 驱动层:金仓官方JDBC驱动成熟稳定,和PG高度兼容
- ORM层:MyBatis-Plus从v3.3.0开始官方支持,开箱即用
- 批处理层:Spring Batch用PostgreSQL模式可以完美运行
对于正在做信创改造的团队来说,Java技术栈迁移到金仓,技术上是完全可行的。关键是要把细节处理好——方言配置、SQL兼容性、连接池调优、批处理参数,每一项都要到位。




