数据开放平台架构设计(四):数据转换与加工
上一篇: 数据开放平台(三):鉴权与访问控制
SQL 查出来的原始数据,通常不能直接给调用方用。数据库字段名是 user_name,调用方要 userName;手机号 13812345678 不能直接暴露,得打成 138****5678;日期格式是 2026-04-23 10:30:00,调用方只要 2026-04-23。
这些"末尾一公里"的数据加工,就是数据转换引擎干的事。它不是替代数仓 ETL,也不应该把复杂计算都塞到接口层;它更适合做轻量、在线、和接口返回强相关的处理。
一、整体设计
数据转换引擎用规则驱动的架构。用户在界面上配置转换规则,运行时引擎按规则依次处理。每条规则是一个独立的 Handler,可插拔、可组合。
数据转换流水线要把字段形态写清楚,避免后续在响应层临时拼逻辑:
@Component
public class DataTransformEngine {
// 所有 TransformHandler 按 Order 排序
@Autowired
private List<TransformHandler> handlers;
public List<Map<String, Object>> transform(TransformConfig config,
List<Map<String, Object>> rawData) {
if (config == null || config.isEmpty()) {
return rawData;
}
List<Map<String, Object>> result = new ArrayList<>(rawData);
for (TransformHandler handler : handlers) {
if (handler.supports(config)) {
result = handler.apply(config, result);
}
}
return result;
}
}生产里我会给转换引擎加两个边界:第一,单次处理行数不能无限大,超过阈值必须分页或异步导出;第二,规则执行要记录耗时,某条自定义正则特别慢时能快速定位。
二、字段映射
更基础的转换:把数据库列名映射为调用方需要的字段名。
自动驼峰转换
数据库字段普遍用下划线命名(user_name),Java/JS 用驼峰命名(userName)。开启自动转换后,全部字段自动处理:
@Component
@Order(1)
public class FieldRenameHandler implements TransformHandler {
@Override
public boolean supports(TransformConfig config) {
return config.getCamelCase() != null || config.getFieldMapping() != null;
}
@Override
public List<Map<String, Object>> apply(TransformConfig config,
List<Map<String, Object>> data) {
return data.stream()
.map(row -> renameFields(config, row))
.collect(Collectors.toList());
}
private Map<String, Object> renameFields(TransformConfig config,
Map<String, Object> row) {
Map<String, Object> result = new LinkedHashMap<>();
for (Map.Entry<String, Object> entry : row.entrySet()) {
String key = entry.getKey();
// 优先用手动映射
String mapped = config.getFieldMapping() != null
? config.getFieldMapping().get(key) : null;
// 没有手动映射,走自动驼峰
if (mapped == null && Boolean.TRUE.equals(config.getCamelCase())) {
mapped = toCamelCase(key);
}
result.put(mapped != null ? mapped : key, entry.getValue());
}
return result;
}
private String toCamelCase(String underscore) {
StringBuilder sb = new StringBuilder();
boolean nextUpper = false;
for (char c : underscore.toCharArray()) {
if (c == '_') {
nextUpper = true;
} else {
sb.append(nextUpper ? Character.toUpperCase(c) : c);
nextUpper = false;
}
}
return sb.toString();
}
}手动字段映射
自动驼峰搞不定的场景(比如 status 要映射为 orderStatus),手动映射覆盖:
{
"fieldMapping": {
"status": "orderStatus",
"create_time": "gmtCreate",
"update_time": "gmtModified"
}
}手动映射优先级高于自动驼峰。映射不到的字段保持原名。
三、数据脱敏
数据开放场景下,脱敏是刚需。用户手机号、身份证号、银行卡号这些敏感信息,必须处理后才能对外暴露。
@Component
@Order(2)
public class DesensitizeHandler implements TransformHandler {
// 内置脱敏规则
private static final Map<String, DesensitizeRule> BUILTIN_RULES = Map.of(
"phone", new RegexRule("(\\d{3})\\d{4}(\\d{4})", "$1****$2"),
"idcard", new RegexRule("(\\d{3})\\d{11}(\\w{4})", "$1***********$2"),
"bankcard", new RegexRule("(\\d{4})\\d+(\\d{4})", "$1 **** **** $2"),
"email", new RegexRule("(\\w)[\\w.]*@(.+)", "$1***@$2"),
"name", new RegexRule("(\\w{1})(.*)", "$1*")
);
@Override
public List<Map<String, Object>> apply(TransformConfig config,
List<Map<String, Object>> data) {
List<DesensitizeField> fields = config.getDesensitizeFields();
if (fields == null || fields.isEmpty()) return data;
return data.stream().map(row -> {
Map<String, Object> result = new LinkedHashMap<>(row);
for (DesensitizeField field : fields) {
Object value = result.get(field.getFieldName());
if (value != null) {
result.put(field.getFieldName(),
desensitize(value.toString(), field.getRule()));
}
}
return result;
}).collect(Collectors.toList());
}
private String desensitize(String value, String rule) {
DesensitizeRule desensitizeRule = BUILTIN_RULES.get(rule);
if (desensitizeRule != null) {
return desensitizeRule.apply(value);
}
// 自定义正则
return value.replaceAll(rule, "***");
}
}配置方式很简单,在界面上勾选字段、选择脱敏规则:
{
"desensitizeFields": [
{ "fieldName": "phone", "rule": "phone" },
{ "fieldName": "id_card", "rule": "idcard" },
{ "fieldName": "email", "rule": "email" }
]
}脱敏后的效果:
13812345678→138****5678110101199001011234→110***********1234zhangsan@example.com→z***@example.com
四、格式转换
数据库里的数据格式和调用方期望的格式经常不一致。日期是更典型的例子,数据库存的是 Timestamp,调用方要的是格式化字符串。
@Component
@Order(3)
public class FormatConvertHandler implements TransformHandler {
@Override
public List<Map<String, Object>> apply(TransformConfig config,
List<Map<String, Object>> data) {
List<FormatRule> rules = config.getFormatRules();
if (rules == null || rules.isEmpty()) return data;
return data.stream().map(row -> {
Map<String, Object> result = new LinkedHashMap<>(row);
for (FormatRule rule : rules) {
Object value = result.get(rule.getFieldName());
if (value != null) {
result.put(rule.getFieldName(),
convert(value, rule.getType(), rule.getPattern()));
}
}
return result;
}).collect(Collectors.toList());
}
private Object convert(Object value, String type, String pattern) {
return switch (type) {
case "date" -> {
LocalDateTime dt = toLocalDateTime(value);
yield dt != null ? dt.format(DateTimeFormatter.ofPattern(
pattern != null ? pattern : "yyyy-MM-dd")) : value;
}
case "datetime" -> {
LocalDateTime dt = toLocalDateTime(value);
yield dt != null ? dt.format(DateTimeFormatter.ofPattern(
pattern != null ? pattern : "yyyy-MM-dd HH:mm:ss")) : value;
}
case "number" -> {
if (value instanceof Number num) {
DecimalFormat df = new DecimalFormat(
pattern != null ? pattern : "#,##0.00");
yield df.format(num.doubleValue());
}
yield value;
}
case "enum" -> {
// 枚举翻译:status=1 → "启用"
if (pattern != null) {
Map<String, String> mapping = JSONUtil.parseObj(pattern);
yield mapping.getOrDefault(value.toString(), value.toString());
}
yield value;
}
default -> value;
};
}
}配置示例:
{
"formatRules": [
{ "fieldName": "create_time", "type": "date", "pattern": "yyyy-MM-dd" },
{ "fieldName": "amount", "type": "number", "pattern": "#,##0.00" },
{ "fieldName": "status", "type": "enum", "pattern": "{\"1\":\"启用\",\"0\":\"禁用\"}" }
]
}枚举翻译特别实用。数据库里 status=1 调用方看不懂,翻译成"启用"就清晰了。配置一次,所有调用方自动受益。
五、字段计算
有时候需要对原始字段做计算,生成新的字段。比如把 first_name 和 last_name 拼成 full_name,或者根据 amount 的范围算出 level。
@Component
@Order(4)
public class FieldComputeHandler implements TransformHandler {
@Override
public List<Map<String, Object>> apply(TransformConfig config,
List<Map<String, Object>> data) {
List<ComputeRule> rules = config.getComputeRules();
if (rules == null || rules.isEmpty()) return data;
return data.stream().map(row -> {
Map<String, Object> result = new LinkedHashMap<>(row);
for (ComputeRule rule : rules) {
Object computed = compute(rule, result);
result.put(rule.getTargetField(), computed);
}
return result;
}).collect(Collectors.toList());
}
private Object compute(ComputeRule rule, Map<String, Object> row) {
return switch (rule.getExpression()) {
case "concat" -> {
// 字段拼接
List<String> fields = rule.getArgs();
String separator = rule.getSeparator() != null
? rule.getSeparator() : "";
yield fields.stream()
.map(f -> String.valueOf(row.getOrDefault(f, "")))
.collect(Collectors.joining(separator));
}
case "substring" -> {
// 字段截取
String value = String.valueOf(row.get(rule.getArgs().get(0)));
int start = Integer.parseInt(rule.getArgs().get(1));
int end = rule.getArgs().size() > 2
? Integer.parseInt(rule.getArgs().get(2)) : value.length();
yield value.substring(Math.min(start, value.length()),
Math.min(end, value.length()));
}
case "condition" -> {
// 条件赋值
Object value = row.get(rule.getArgs().get(0));
Map<String, String> mapping = JSONUtil.parseObj(rule.getPattern());
yield mapping.getOrDefault(
value != null ? value.toString() : "null",
rule.getDefaultValue());
}
default -> null;
};
}
}配置示例,拼接姓名:
{
"computeRules": [
{
"targetField": "fullName",
"expression": "concat",
"args": ["first_name", "last_name"],
"separator": " "
}
]
}六、聚合计算
有些场景需要对结果集做聚合,总数、平均值、上限值、分组统计。比如"每个部门的平均薪资"、“每月的订单总额”。
@Component
@Order(5)
public class AggregateHandler implements TransformHandler {
@Override
public boolean supports(TransformConfig config) {
return config.getAggregateConfig() != null;
}
@Override
public List<Map<String, Object>> apply(TransformConfig config,
List<Map<String, Object>> data) {
AggregateConfig agg = config.getAggregateConfig();
if ("groupBy".equals(agg.getType())) {
return groupBy(data, agg.getGroupField(), agg.getMetrics());
}
if ("summary".equals(agg.getType())) {
return Collections.singletonList(
computeSummary(data, agg.getMetrics()));
}
return data;
}
private List<Map<String, Object>> groupBy(List<Map<String, Object>> data,
String groupField,
List<MetricConfig> metrics) {
Map<Object, List<Map<String, Object>>> groups = data.stream()
.collect(Collectors.groupingBy(row -> row.get(groupField)));
return groups.entrySet().stream().map(entry -> {
Map<String, Object> result = new LinkedHashMap<>();
result.put(groupField, entry.getKey());
for (MetricConfig metric : metrics) {
result.put(metric.getAlias(),
computeMetric(entry.getValue(), metric));
}
return result;
}).collect(Collectors.toList());
}
private Object computeMetric(List<Map<String, Object>> rows,
MetricConfig metric) {
return switch (metric.getFunction()) {
case "count" -> rows.size();
case "sum" -> rows.stream()
.map(row -> toBigDecimal(row.get(metric.getField())))
.filter(Objects::nonNull)
.reduce(BigDecimal.ZERO, BigDecimal::add);
case "avg" -> {
BigDecimal sum = rows.stream()
.map(row -> toBigDecimal(row.get(metric.getField())))
.filter(Objects::nonNull)
.reduce(BigDecimal.ZERO, BigDecimal::add);
yield rows.isEmpty() ? BigDecimal.ZERO
: sum.divide(BigDecimal.valueOf(rows.size()),
2, RoundingMode.HALF_UP);
}
case "max" -> rows.stream()
.map(row -> toBigDecimal(row.get(metric.getField())))
.filter(Objects::nonNull)
.max(BigDecimal::compareTo)
.orElse(null);
case "min" -> rows.stream()
.map(row -> toBigDecimal(row.get(metric.getField())))
.filter(Objects::nonNull)
.min(BigDecimal::compareTo)
.orElse(null);
default -> null;
};
}
}聚合计算配置示例:
{
"aggregateConfig": {
"type": "groupBy",
"groupField": "department",
"metrics": [
{ "field": "salary", "function": "avg", "alias": "avgSalary" },
{ "field": "id", "function": "count", "alias": "headcount" }
]
}
}聚合这件事要克制。能在 SQL 里完成的聚合,优先交给数据库;平台侧聚合适合做小结果集的补充口径,比如接口已经查出了几十条明细,再顺手算一个 summary。不要把它当成实时数仓来用。
七、规则组合与执行顺序
多个转换规则可以自由组合。执行顺序由 @Order 注解控制:
- 字段重命名(先统一字段名)
- 数据脱敏(脱敏后再做其他处理)
- 格式转换(格式化输出)
- 字段计算(基于已有字段生成新字段)
- 聚合计算(末尾做聚合,输入已经是加工过的数据)
这个顺序是有讲究的。字段映射要尽量靠前,后面的规则才能用统一字段名;脱敏要在响应前完成,避免中间结果被直接返回;字段计算要在聚合之前,先算出人均消费,再按部门聚合。
如果默认顺序不满足需求,用户可以自定义执行顺序。界面上拖拽调整规则顺序,保存后运行时按新顺序执行。
上一篇: 数据开放平台(三):鉴权与访问控制
下一篇: 数据开放平台(五):监控与治理