← 返回

数据开放平台架构设计(四):数据转换与加工

上一篇: 数据开放平台(三):鉴权与访问控制

SQL 查出来的原始数据,通常不能直接给调用方用。数据库字段名是 user_name,调用方要 userName;手机号 13812345678 不能直接暴露,得打成 138****5678;日期格式是 2026-04-23 10:30:00,调用方只要 2026-04-23

这些"末尾一公里"的数据加工,就是数据转换引擎干的事。它不是替代数仓 ETL,也不应该把复杂计算都塞到接口层;它更适合做轻量、在线、和接口返回强相关的处理。


一、整体设计

数据转换引擎用规则驱动的架构。用户在界面上配置转换规则,运行时引擎按规则依次处理。每条规则是一个独立的 Handler,可插拔、可组合。

数据转换流水线要把字段形态写清楚,避免后续在响应层临时拼逻辑:

flowchart LR A["原始结果集<br/>db_column"] --> B["字段映射<br/>apiField"] B --> C["脱敏<br/>maskedField"] C --> D["格式化<br/>displayValue"] D --> E["字段计算<br/>computedField"] E --> F["小结果集聚合<br/>summary"] F --> G["字段权限过滤<br/>visibleFields"] G --> H["响应包装<br/>data、meta、traceId"]
@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" }
  ]
}

脱敏后的效果:

  • 13812345678138****5678
  • 110101199001011234110***********1234
  • zhangsan@example.comz***@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_namelast_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 注解控制:

  1. 字段重命名(先统一字段名)
  2. 数据脱敏(脱敏后再做其他处理)
  3. 格式转换(格式化输出)
  4. 字段计算(基于已有字段生成新字段)
  5. 聚合计算(末尾做聚合,输入已经是加工过的数据)

这个顺序是有讲究的。字段映射要尽量靠前,后面的规则才能用统一字段名;脱敏要在响应前完成,避免中间结果被直接返回;字段计算要在聚合之前,先算出人均消费,再按部门聚合。

如果默认顺序不满足需求,用户可以自定义执行顺序。界面上拖拽调整规则顺序,保存后运行时按新顺序执行。


上一篇: 数据开放平台(三):鉴权与访问控制
下一篇: 数据开放平台(五):监控与治理