✅ 一套 基于Java 17+的数据质量规则DSL设计(告别硬编码,规则配置化)。
✅ 一个 金仓方言SQL生成器(完美适配KingbaseES的正则、日期、空串怪癖)。
✅ 一套 生产级分片执行引擎(基于主键分片+JDBC流式读取,千万级大表校验不OOM、不锁表!)。
✅ 人大金仓数据校验的 “三大暗坑”避坑指南。
收藏这篇,下次信创数据迁移前跑一遍,能让你少背几个P0级故障的锅,安稳睡个好觉。
一、为什么传统的数据校验在“人大金仓”上跑不通?
在动手写代码前,必须先搞清楚我们在跟什么“怪物”搏斗。很多老铁觉得:“数据校验嘛,不就是写几个SQL查一下 WHERE column IS NULL 吗?”
Too young too simple! 在信创环境下,传统校验方式有“三大死穴”。
死穴1:大表全表扫描,直接把金仓“扫宕机”
政务中台的核心表(如 pop_base_info 人口基础信息表)动辄几千万甚至上亿行。如果你直接跑 SELECT COUNT(*) FROM pop_base_info WHERE id_card IS NULL,金仓会走全表扫描(Seq Scan)。在业务高峰期,这种大查询会瞬间吃满CPU和IO,导致正常的OLTP事务被阻塞,直接引发生产事故。
死穴2:方言差异,让你的校验SQL“集体失效”
你从MySQL/Oracle抄来的校验SQL,在金仓里可能直接报语法错误:
正则校验:Oracle用 REGEXP_LIKE,MySQL用 REGEXP,金仓(PG系)用 ~ 或 SIMILAR TO。
空串与NULL:金仓Oracle兼容模式下 ‘’ = NULL,你写的 WHERE remark = ‘’ 永远查不出数据!
日期校验:金仓对日期格式极其严格,没有 ISDATE() 这种傻瓜函数,必须自己写正则或 TO_DATE 捕获异常。
死穴3:结果集OOM,Java应用直接“暴毙”
如果你用MyBatis或JDBC把“不合规的脏数据”查出来展示,*千万不要 SELECT!如果有100万条脏数据,JDBC驱动会试图把它们全加载到JVM内存里,直接触发 java.lang.OutOfMemoryError: Java heap space。
二、破局之道:轻量级数据质量监控引擎(DQM Engine)架构
怎么破局?我的思路是:规则DSL化 + 方言适配 + 分片流式执行。
graph TD
subgraph 规则配置层
R1[YAML/JSON 规则定义] --> R2[规则解析器]
end
subgraph 引擎核心层 (DQM Engine) R2 --> E1[金仓方言SQL生成器] E1 --> E2[分片调度器 ShardScheduler] E2 --> E3[JDBC流式读取器] end subgraph 执行与告警层 E3 --> DB[(人大金仓 KingbaseES)] E3 --> M1[指标聚合器] M1 --> A1[飞书/钉钉/邮件 告警] M1 --> P1[质量报告持久化] end style E1 fill:#ff6600,color:#fff style E2 fill:#4488ff,color:#fff设计核心思想:
规则与代码解耦:用YAML定义校验规则(如:非空、正则、枚举、外键),引擎负责解析。
方言隔离:SQL生成器根据数据库类型(Kingbase/DM/OB)动态生成对应的方言SQL。
分片与流式:把千万级大表按主键(ID)切分成多个小分片(如每次10万条),并使用JDBC的 setFetchSize 流式读取,内存占用始终控制在MB级别!
三、核心代码实现(生产级,极度详尽)
老铁们,泡好咖啡,下面这几百行Java 17代码是我熬了无数个通宵打磨出来的。每一行注释都是真金白银的踩坑经验,涵盖逻辑、边界、性能与易错点。
3.1 模块一:数据质量规则模型与DSL设计
我们使用Java 17的 Record 和 Sealed Interface 来构建不可变的规则模型,确保线程安全和类型安全。
package com.mobai.dqm.model;
import java.util.List;
/**
🔧 模块名称: QualityRule (数据质量规则模型)
📝 功能描述: 使用 Java 17 Sealed Interface 定义封闭的规则类型体系
🏗️ 设计思想:
规则不可变(Record),天然线程安全,适合在并发调度引擎中传递。
封闭接口(Sealed),限制规则类型,便于后续使用 Switch Expression 进行模式匹配。
*/
public sealed interface QualityRule permits
QualityRule.NotNull,
QualityRule.RegexMatch,
QualityRule.LengthRange,
QualityRule.EnumIn,
QualityRule.CustomSql {
String columnName(); String ruleName(); SeverityLevel severity(); // 严重级别:BLOCK(阻断), WARN(警告), INFO(提示) /** 1. 非空校验(注意:在金仓中需要特殊处理空字符串) */ record NotNull(String columnName, String ruleName, SeverityLevel severity, boolean treatEmptyStringAsNull) implements QualityRule {} /** 2. 正则校验(如身份证、手机号) */ record RegexMatch(String columnName, String ruleName, SeverityLevel severity, String regexPattern) implements QualityRule {} /** 3. 长度范围校验 */ record LengthRange(String columnName, String ruleName, SeverityLevel severity, int min, int max) implements QualityRule {} /** 4. 枚举值校验 */ record EnumIn(String columnName, String ruleName, SeverityLevel severity, List<String> allowedValues) implements QualityRule {} /** 5. 自定义SQL校验(兜底方案,返回违规数量) */ record CustomSql(String columnName, String ruleName, SeverityLevel severity, String sqlTemplate) implements QualityRule {}}
enum SeverityLevel {
BLOCK, WARN, INFO
}
对应的 YAML 配置文件示例 (rules/pop_base_info.yml):
table: pop_base_info
primary_key: id
rules:
column: id_card
type: RegexMatch
name: “身份证号格式校验”
severity: BLOCK
pattern: “1\d{5}(18|19|20)\d{2}((0[1-9])|(1[0-2]))(([0-2][1-9])|10|20|30|31)\d{3}[0-9Xx]”
column: phone
type: NotNull
name: “手机号非空”
severity: WARN
treat_empty_string_as_null: true # ⚠️ 金仓专属配置:把’'也当成NULL来查
3.2 模块二:金仓方言SQL生成器(核心!)
这是整个引擎的“翻译官”。它把抽象的规则,翻译成人大金仓(KingbaseES)能高效执行的SQL。
package com.mobai.dqm.dialect;
import com.mobai.dqm.model.QualityRule;
import org.springframework.stereotype.Component;
import java.util.stream.Collectors;
/**
🔧 模块名称: KingbaseSqlDialectGenerator
📝 功能描述: 人大金仓(KingbaseES)专属 SQL 生成器
🏗️ 设计思想:
利用 Java 17 的 Switch Pattern Matching,将规则对象精准映射为金仓方言 SQL。
⚠️ 易错点:
金仓 Oracle 模式下 ‘’ = NULL,非空校验必须加上 OR column = ‘’。
金仓的正则匹配操作符是 ‘~’ (区分大小写) 或 ‘~*’ (不区分),不是 REGEXP_LIKE。
字符串长度计算:char_length 按字符算,octet_length 按字节算。业务通常用 char_length。
*/
@Component
public class KingbaseSqlDialectGenerator {
/** 生成统计“违规数据数量”的 SQL * @param tableName 表名 @param rule 质量规则 @param shardCondition 分片条件(如:id >= 1 AND id < 10000),防止全表扫描 @return 可执行的 SQL */ public String generateCountSql(String tableName, QualityRule rule, String shardCondition) { String whereClause = buildWhereClause(rule); // 💡 性能考量:强制带上分片条件,将大查询拆解为小查询 // 使用 COUNT) 而不是 COUNT(1) 或 COUNT(col),在 PG/金仓 底层优化器中 COUNT() 最快 return String.format(""" SELECT COUNT(*) FROM %s WHERE (%s) AND (%s) """, tableName, shardCondition, whereClause); } /** 生成采样“违规数据明细”的 SQL(用于告警展示,限制条数) */ public String generateSampleSql(String tableName, QualityRule rule, String shardCondition, int limit) { String whereClause = buildWhereClause(rule); return String.format(""" SELECT * FROM %s WHERE (%s) AND (%s) LIMIT %d """, tableName, shardCondition, whereClause, limit); } /** 核心:根据规则类型构建 WHERE 条件 💡 使用 Java 17 Switch Pattern Matching,代码极其优雅且编译器保证穷举 */ private String buildWhereClause(QualityRule rule) { return switch (rule) { case QualityRule.NotNull r -> { // 🚨 金仓大坑:如果 treatEmptyStringAsNull 为 true,必须显式加上 OR col = '' // 因为在金仓里,'' 可能被底层转成了 NULL,但也可能在某些兼容模式下保留 // 双管齐下,确保脏数据无处遁形 if (r.treatEmptyStringAsNull()) { yield String.format("(%s IS NULL OR %s = '')", r.columnName(), r.columnName()); } else { yield String.format("%s IS NULL", r.columnName()); } } case QualityRule.RegexMatch r -> { // ⚠️ 金仓/PG 的正则操作符是 '~' (匹配) // 我们要查的是“不合规”的数据,所以用 '!~' (不匹配) // 同时必须排除 NULL 值,否则 NULL !~ 'regex' 的结果是 NULL,不会被 WHERE 过滤 yield String.format("(%s IS NOT NULL AND %s !~ '%s')", r.columnName(), r.columnName(), escapeSql(r.regexPattern())); } case QualityRule.LengthRange r -> { // 使用 char_length 计算字符数(兼容中文) yield String.format("(char_length(%s) < %d OR char_length(%s) > %d)", r.columnName(), r.min(), r.columnName(), r.max()); } case QualityRule.EnumIn r -> { String inValues = r.allowedValues().stream() .map(v -> "'" + escapeSql(v) + "'") .collect(Collectors.joining(",")); // 不在枚举列表中,且排除 NULL(如果 NULL 是允许的,需另行配置) yield String.format("(%s IS NOT NULL AND %s NOT IN (%s))", r.columnName(), r.columnName(), inValues); } case QualityRule.CustomSql r -> { // 自定义 SQL 直接透传,但要求 DBA 保证 SQL 的安全性 yield r.sqlTemplate(); } }; } /** 简单的 SQL 注入防护(转义单引号) ⚠️ 边界处理:防止规则配置中的单引号破坏 SQL 语法 */ private String escapeSql(String input) { if (input == null) return ""; return input.replace("'", "''"); }}
3.3 模块三:分布式分片执行引擎(解决大表OOM的终极杀器)
这是整个系统最硬核的部分。面对千万级大表,我们绝对不能一条SQL扫到底。必须按主键分片,并使用JDBC流式读取。
package com.mobai.dqm.engine;
import com.mobai.dqm.dialect.KingbaseSqlDialectGenerator;
import com.mobai.dqm.model.QualityRule;
import com.mobai.dqm.model.SeverityLevel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
/**
🔧 模块名称: ShardExecutionEngine
📝 功能描述: 基于主键分片 + JDBC流式读取的数据质量执行引擎
🏗️ 设计思想:
获取表的主键最大值和最小值,按固定步长(如10万)切分为多个分片。
使用线程池并发执行各个分片的校验SQL。
采样明细时,强制使用 JDBC 的 setFetchSize 开启游标流式读取,防止 OOM。
⚡ 性能目标: 千万级表的全量规则校验 < 5分钟,内存占用 < 500MB。
*/
@Component
public class ShardExecutionEngine {
private static final Logger log = LoggerFactory.getLogger(ShardExecutionEngine.class); // 💡 分片大小:每次扫描 10万 条数据。 // 太小会导致分片过多,线程切换和SQL解析开销大;太大会导致单次查询时间过长,容易超时。 private static final int SHARD_SIZE = 100_000; private final JdbcTemplate jdbcTemplate; private final KingbaseSqlDialectGenerator dialectGenerator; // ⚡ 性能优化:使用自定义线程池,核心线程数根据 CPU 核心数动态调整 // 拒绝策略使用 CallerRunsPolicy,防止任务堆积导致 OOM,宁可让主线程慢点跑 private final ExecutorService executorService = new ThreadPoolExecutor( Runtime.getRuntime().availableProcessors(), Runtime.getRuntime().availableProcessors() * 2, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(1000), new ThreadPoolExecutor.CallerRunsPolicy() ); public ShardExecutionEngine(JdbcTemplate jdbcTemplate, KingbaseSqlDialectGenerator dialectGenerator) { this.jdbcTemplate = jdbcTemplate; this.dialectGenerator = dialectGenerator; } /** 执行单张表的单条规则校验 * @return 违规数据总数 */ public long executeRule(String tableName, String primaryKey, QualityRule rule) { log.info("🚀 开始校验: 表={}, 规则={}", tableName, rule.ruleName()); // 1. 获取主键的 Min 和 Max 值,确定分片边界 // ⚠️ 易错点:如果表是空的,MAX() 会返回 NULL,必须做防空处理 Long minId = jdbcTemplate.queryForObject( String.format("SELECT MIN(%s) FROM %s", primaryKey, tableName), Long.class); Long maxId = jdbcTemplate.queryForObject( String.format("SELECT MAX(%s) FROM %s", primaryKey, tableName), Long.class); if (minId == null || maxId == null) { log.warn("表 {} 为空,跳过校验。", tableName); return 0; } // 2. 切分分片任务 long totalViolations = 0; long currentStart = minId; CompletableFuture<Long>[] futures = new CompletableFuture[(int) ((maxId - minId) / SHARD_SIZE) + 1]; int futureIndex = 0; while (currentStart <= maxId) { long currentEnd = Math.min(currentStart + SHARD_SIZE - 1, maxId); String shardCondition = String.format("%s >= %d AND %s <= %d", primaryKey, currentStart, primaryKey, currentEnd); // 生成该分片的统计 SQL String countSql = dialectGenerator.generateCountSql(tableName, rule, shardCondition); // 🚀 异步提交分片任务 futures[futureIndex++] = CompletableFuture.supplyAsync(() -> { try { Long count = jdbcTemplate.queryForObject(countSql, Long.class); return count != null ? count : 0L; } catch (Exception e) { log.error("❌ 分片执行失败: SQL={}, 原因={}", countSql, e.getMessage()); return 0L; // 降级处理:单个分片失败不影响全局,记0并告警 } }, executorService); currentStart += SHARD_SIZE; } // 3. 等待所有分片完成,并汇总结果 try { CompletableFuture.allOf(futures).get(10, TimeUnit.MINUTES); // 设置全局超时时间 for (int i = 0; i < futureIndex; i++) { totalViolations += futures[i].get(); } } catch (Exception e) { log.error("❌ 规则执行超时或异常: {}", rule.ruleName(), e); } log.info("✅ 校验完成: 表={}, 规则={}, 违规数={}", tableName, rule.ruleName(), totalViolations); return totalViolations; } /** 🚨 核心防OOM设计:流式采样违规数据明细 * 💡 为什么不用 JdbcTemplate? 因为 JdbcTemplate 默认会把 ResultSet 全部加载到内存。 我们必须拿到原生的 Connection 和 Statement,设置 setFetchSize, 强制 JDBC 驱动使用“游标(Cursor)”模式,每次只从金仓拉取 100 条数据! */ public void streamSampleViolations(String tableName, String primaryKey, QualityRule rule, DataSource dataSource, ViolationConsumer consumer) { String sampleSql = dialectGenerator.generateSampleSql(tableName, rule, "1=1", 10000); // ⚠️ 金仓/PG 开启游标流式读取的 3 个必要条件: // 1. autocommit 必须为 false // 2. setFetchSize 必须 > 0 // 3. ResultSet 类型必须是 FORWARD_ONLY try (Connection conn = dataSource.getConnection()) { conn.setAutoCommit(false); try (Statement stmt = conn.createStatement(ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY)) { stmt.setFetchSize(100); // 每次从网络拉取 100 行 try (ResultSet rs = stmt.executeQuery(sampleSql)) { int metaCount = rs.getMetaData().getColumnCount(); while (rs.next()) { // 将行数据封装为 Map 传给消费者(如写入文件或发送告警) // 这里绝对不能把 rs 对象传出去,因为游标随时会关闭! java.util.Map<String, Object> rowData = new java.util.HashMap<>(); for (int i = 1; i <= metaCount; i++) { rowData.put(rs.getMetaData().getColumnName(i), rs.getObject(i)); } consumer.accept(rowData); } } } } catch (Exception e) { log.error("❌ 流式采样失败: {}", e.getMessage(), e); } } @FunctionalInterface public interface ViolationConsumer { void accept(java.util.Map<String, Object> rowData); }}
💡 工程实践总结(防OOM三连):
分片(Sharding):把 WHERE id BETWEEN 1 AND 100000 拆解,让金仓走主键索引(Index Scan),避免全表扫描(Seq Scan)导致的IO打满。
游标(Cursor):setFetchSize(100) + setAutoCommit(false) 是 PG/金仓 驱动开启流式读取的“咒语”。少一个,JVM都会因为加载百万行数据而OOM。
超时控制(Timeout):CompletableFuture.allOf().get(timeout) 防止某个分片因为锁等待或慢查询卡死整个线程池。
四、人大金仓数据校验的“三大暗坑”避坑指南
代码写好了,但如果你不懂金仓底层的“怪脾气”,跑起来照样翻车。这是我用无数个“回滚”和“告警”换来的血泪教训。
🔴 暗坑1:ora_input_emptystr_isnull 导致的“非空校验”漏网之鱼
现象:你在Java里配置了 NotNull 规则,生成的SQL是 WHERE phone IS NULL。结果跑出来违规数是0,但业务层还是报NPE!
根因:金仓在Oracle兼容模式下,参数 ora_input_emptystr_isnull = on。业务代码插入了空字符串 ‘’,金仓底层把它存成了 NULL。但是,某些老数据是通过DataX或Kettle直接绕过应用层写入的,这些工具可能保留了物理上的 ‘’(空串)。
避坑:在生成非空校验SQL时,必须写成 WHERE col IS NULL OR col = ‘’(如上文代码所示)。双管齐下,绝不漏掉任何一条脏数据。
🟠 暗坑2:正则校验中的“转义字符”地狱
现象:校验手机号的正则 ^1[3-9]d{9},在Java字符串里你要写成 ^1[3-9]\d{9}。当你把它拼接到金仓的SQL里时,金仓的解析器可能会把 d 当成普通的 d,导致正则完全失效,把所有数据都判定为“违规”!
避坑:在金仓(PG系)中,如果使用标准正则操作符 ~,必须使用“转义字符串”语法(E’…‘)。
正确写法:WHERE phone !~ E’^1[3-9]\d{9}'。在上面的 KingbaseSqlDialectGenerator 中,为了简化,我使用了 \d,但在实际生产环境中,建议在拼接SQL时显式加上 E 前缀,或者使用 d 的等价类 [0-9] 来彻底规避转义问题(如代码中所示,我直接用了 \d,在Java 15+ 的 Text Block 中处理会更优雅)。
🟡 暗坑3:系统表查询的“权限隔离”
现象:你想写个规则,校验“表是否存在”或者“索引是否失效”,去查 pg_class 或 sys_tables。结果用业务账号一连,直接报 permission denied for table pg_class。
避坑:金仓(继承PG)对系统目录的权限控制极其严格。数据质量监控引擎必须使用独立的“监控专属账号”,并且要由DBA授予 pg_monitor 角色或特定的 SELECT 权限。千万别用业务账号(如 app_user)去跑监控SQL!
五、总结与互动
金句总结
💡 “没有自动化监控的数据治理,就是‘盲人摸象’。你以为数据很干净,其实只是脏数据还没触发业务报错。”
💡 “在千万级大表面前,任何不带‘分片’和‘流式’的查询,都是对数据库IO的恐怖袭击,也是对JVM内存的蓄意谋杀。”
💡 “信创迁移不是简单的‘数据搬家’,而是借机用‘数据质量引擎’给历史屎山做一次彻底的‘肠胃镜’。”
本文知识点回顾
mindmap
root((Java 数据质量监控引擎))
核心痛点
大表全扫OOM
金仓方言不兼容
规则硬编码难维护
架构设计
YAML 规则 DSL
金仓方言 SQL 生成器
分片与流式执行引擎
工程落地
Java 17 Record/Sealed
主键分片并发查询
JDBC FetchSize 游标读取
金仓避坑
空串与NULL的罗生门
正则转义字符地狱
系统表权限隔离
1-9 ↩︎