# 数据源头：业务数据与用户日志双生成器

> 讲解仿真数据是怎么造出来的：一个 Spring 应用把商品、订单、用户等全链路业务数据写入 MySQL，另一个 Spring 应用持续构造带嵌套结构的用户行为日志并发往 Kafka，为湖仓提供真实感输入。

- Repository: fengyu-eng/paimon-datalake
- GitHub: https://github.com/fengyu-eng/paimon-datalake
- Human wiki: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533
- Complete Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/llms-full.txt

## Source Files

- `src/main/datagenerator/business_code/BusinessApplication.java`
- `src/main/datagenerator/business_code/generator/BaseDataGenerator.java`
- `src/main/datagenerator/userlog_code/UserLogApplication.java`
- `src/main/datagenerator/userlog_code/generator/UserLogGenerator.java`
- `src/main/datagenerator/business_code/util/DbUtil.java`
- `src/main/datagenerator/business_code/util/RandomUtil.java`

---

<details>
<summary>相关源文件</summary>
以下文件用于生成此维基页面：
- [src/main/datagenerator/business_code/BusinessApplication.java](src/main/datagenerator/business_code/BusinessApplication.java)
- [src/main/datagenerator/business_code/generator/BaseDataGenerator.java](src/main/datagenerator/business_code/generator/BaseDataGenerator.java)
- [src/main/datagenerator/business_code/generator/UserDataGenerator.java](src/main/datagenerator/business_code/generator/UserDataGenerator.java)
- [src/main/datagenerator/business_code/generator/OrderDataGenerator.java](src/main/datagenerator/business_code/generator/OrderDataGenerator.java)
- [src/main/datagenerator/business_code/generator/UserBehaviorGenerator.java](src/main/datagenerator/business_code/generator/UserBehaviorGenerator.java)
- [src/main/datagenerator/business_code/util/DbUtil.java](src/main/datagenerator/business_code/util/DbUtil.java)
- [src/main/datagenerator/business_code/util/RandomUtil.java](src/main/datagenerator/business_code/util/RandomUtil.java)
- [src/main/datagenerator/userlog_code/UserLogApplication.java](src/main/datagenerator/userlog_code/UserLogApplication.java)
- [src/main/datagenerator/userlog_code/generator/UserLogGenerator.java](src/main/datagenerator/userlog_code/generator/UserLogGenerator.java)
- [src/main/datagenerator/userlog_code/model/UserLog.java](src/main/datagenerator/userlog_code/model/UserLog.java)
- [src/main/datagenerator/userlog_code/model/Common.java](src/main/datagenerator/userlog_code/model/Common.java)
</details>

# 数据源头：业务数据与用户日志双生成器

## 页面概述

这个仓库做的是"从 0 到 1 构建电商实时数据湖仓"，整个链路需要有源源不断的输入数据。本页讲解的就是这套输入数据是怎么"造"出来的：仓库里放了两套互相独立的 Spring 程序——**业务数据生成器**负责把商品、订单、用户等全链路电商业务数据批量写进 MySQL，**用户日志生成器**负责按固定节奏构造带嵌套结构的用户行为日志并发往 Kafka。

之所以拆成两条线，是因为湖仓对两类数据的采集方式天然不同：业务表走 Flink CDC（读 MySQL 的变更），用户日志走 Kafka 流。两个生成器正好把两类"有真实感"的输入喂到下游，是整条 ODS 层数据的第一来源。

Sources: [README.md:6-10]()；[src/main/datagenerator/userlog_code/README.md:1-4]()；[src/main/datagenerator/business_code/README.md:1-4]()

## 两个程序在整个湖仓里的位置

可以把数据源头理解为一条"双叉路口"：左边是结构化业务表，右边是半结构化日志流，汇合点在下游的 ODS 层。

```mermaid
flowchart LR
    subgraph SRC["数据源头（src/main/datagenerator 下的 Spring 进程）"]
        BA["BusinessApplication<br/>业务数据生成器"]
        UA["UserLogApplication<br/>用户日志生成器"]
    end

    subgraph CARRIER["数据载体"]
        DB[("MySQL<br/>gmall 库")]
        KAFKA["Kafka 主题<br/>ODS_BASE_LOG"]
    end

    subgraph LAKE["湖仓采集（Flink SQL）"]
        CDC["Flink CDC<br/>mysql-cdc 连接器"]
        KC["Flink Kafka 连接器"]
        ODS["Paimon ODS 层"]
    end

    BA -- "HikariCP 批量 INSERT" --> DB
    UA -- "逐条 JSON 消息" --> KAFKA
    DB --> CDC --> ODS
    KAFKA --> KC --> ODS
```

Sources: [src/main/datagenerator/business_code/BusinessApplication.java:13-78]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:22-100]()；[src/main/java/com/fy/warehouse/ODS/ods_log_inc.java:24-45]()；[src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java:21-47]()

一个容易混淆的点先说明：`business_jar/` 与 `userlog_jar/` 目录里放的是另一套预构建 mock 工具的部署脚本与配置（`log.sh` 启动 `datageneration2024-mock-log-2024-05-01.jar`，配置键以 `mock.*` 开头），仓库内并没有这个 jar 的本体。本页描述的是仓库内带源码的 `*_code` 两套实现。

Sources: [src/main/datagenerator/business_jar/log.sh:1-4]()；[src/main/datagenerator/userlog_jar/application.yml:1-12]()；[src/main/datagenerator/business_jar/application.properties:1-9]()

## 两个生成器的整体对照

| 对比维度 | 业务数据生成器 | 用户日志生成器 |
| --- | --- | --- |
| 入口类 | `BusinessApplication`（business_code） | `UserLogApplication`（userlog_code） |
| 产物形态 | 一张张关系表里的行 | 一条条嵌套 JSON 日志 |
| 数据出口 | MySQL `gmall` 库（JDBC 批量插入） | Kafka 主题（默认 `ODS_BASE_LOG`） |
| 生成节奏 | 主循环每 `interval` 毫秒跑一轮，每轮写一批（默认 1000） | 每 `interval` 毫秒生成并同步发送一条 |
| 关键配置 | `spring.datasource.*`、`generator.batch-size`、`generator.interval` | `kafka.bootstrap-servers`、`kafka.topic`、`generator.interval` |
| 连不上目标时的降级 | 把 SQL 打印到控制台 | 把日志 JSON 打印到控制台 |
| 下游消费方式 | Flink CDC 读 MySQL 变更 | Flink Kafka Connector 读主题 |

Sources: [src/main/datagenerator/business_code/BusinessApplication.java:17-21]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:27-51]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:61-89]()

## 业务数据生成器：把全链路业务表写进 MySQL

### 启动与主循环

`BusinessApplication` 是标准的 Spring Boot 应用（`@SpringBootApplication`），启动后先做一次**基础数据**生成，然后进入无限循环，每一轮依次调用各类生成器，最后 `Thread.sleep(interval)` 控制节奏（默认 5 秒一轮）：

```java
@Override
public void run(String... args) throws Exception {
    baseDataGenerator.generateBaseData(batchSize);      // 一次性：分类/品牌/属性/地区/字典

    while (true) {
        userDataGenerator.generateUserData(batchSize / 10);
        productDataGenerator.generateProductData(batchSize / 10, batchSize / 20);
        activityDataGenerator.generateActivityData(batchSize / 10, batchSize / 20);
        couponDataGenerator.generateCouponData(batchSize / 10, batchSize / 20);
        orderDataGenerator.generateOrderData(batchSize);       // 每轮量最大的交易数据
        userBehaviorGenerator.generateUserBehaviorData(batchSize);
        warehouseDataGenerator.generateWarehouseData(batchSize / 5);
        cmsDataGenerator.generateCMSData(batchSize / 20, batchSize / 40, batchSize / 100);
        Thread.sleep(interval);
    }
}
```

Sources: [src/main/datagenerator/business_code/BusinessApplication.java:57-78]()

### 每轮覆盖哪些数据

在默认 `batch-size=1000` 下，每个循环周期各类生成器会新增大约下面这些行（各子表数量基本随主参数等比缩放）：

| 生成器 | 覆盖的数据/表 | 每轮新增量（默认值） |
| --- | --- | --- |
| `BaseDataGenerator` | 三级分类、品牌、商品属性、销售属性、省份地区、字典、前端参数 | 启动时一次性生成 |
| `UserDataGenerator` | `user_info`、`user_address` | 各约 100 行 |
| `ProductDataGenerator` | `spu_info`、`sku_info` 及其图片/属性/销售属性 | 约 100 个 SPU + 50 个 SKU |
| `ActivityDataGenerator` | `activity_info`、活动规则、活动商品、秒杀 | 约 100 个活动 |
| `CouponDataGenerator` | `coupon_info`、券范围、领券记录 | 约 100 张券 |
| `OrderDataGenerator` | 订单主表 + 明细 + 支付 + 状态流水 + 优惠分摊 + 退款 | 每个子表约 1000 行 |
| `UserBehaviorGenerator` | `cart_info` 购物车、`comment_info` 评价、`favor_info` 收藏 | 每个子表约 1000 行 |
| `WarehouseDataGenerator` / `CMSDataGenerator` | 仓储库存 / Banner、主题、评论 | 相对少量 |

其中 `OrderDataGenerator` 内部又串了订单主表、订单明细、支付信息、订单状态流水、明细-活动、明细-优惠券、退款单、退款流水 8 张表；金额上会用"原价 − 活动优惠 − 券优惠 + 运费"拼出 `total_amount`，让数值之间自洽。

Sources: [src/main/datagenerator/business_code/generator/OrderDataGenerator.java:22-31]()、[OrderDataGenerator.java:49-56]()；[src/main/datagenerator/business_code/generator/UserBehaviorGenerator.java:24-28]()；[src/main/datagenerator/business_code/README.md:6-45]()、[README.md:91-105]()

### 主键续接与"随机外键"

每个生成方法开头都会先查这张表的当前最大主键，再接着往下排：

```java
String maxIdSql = "SELECT COALESCE(MAX(id), 0) FROM base_category1";
int startId = dbUtil.queryForInt(maxIdSql) + 1;
```

Sources: [src/main/datagenerator/business_code/generator/BaseDataGenerator.java:38-41]()；[src/main/datagenerator/business_code/generator/UserDataGenerator.java:27-30]()；[src/main/datagenerator/business_code/util/DbUtil.java:161-178]()

这样做的效果是：**主键单调递增、进程重启不冲突**，配合 CDC 增量同步时能持续产出"新行"。但要特别说明：表之间的关联（订单→用户、明细→SKU、地址→省份）不是真正去查已存在的主数据，而是用随机数落进一个区间，例如 `user_id = RandomUtil.generateNumber(1, 1000)`、省份 ID 取 1–34。所以这是一种"概率意义上的弱外键"——数据形态看起来真实，但同一轮生成的行之间并不保证严格可 join。仓库 README 中"SKU 关联 SPU、订单关联用户"等说法，本质是靠这个随机区间近似实现的。

Sources: [src/main/datagenerator/business_code/generator/OrderDataGenerator.java:63]()、[OrderDataGenerator.java:108-109]()、[OrderDataGenerator.java:139]()；[src/main/datagenerator/business_code/generator/UserDataGenerator.java:73-76]()；[src/main/datagenerator/business_code/README.md:100-105]()

### DbUtil：连接池、批量写与"无库可写"降级

`DbUtil` 用 HikariCP 维护一个连接池，连接串、账号、池大小等全部来自 `spring.datasource.*` 配置。批量写通过 `PreparedStatement.addBatch()` + `executeBatch()` 一次提交：

Sources: [src/main/datagenerator/business_code/util/DbUtil.java:60-71]()、[DbUtil.java:88-113]()

最有意思的是它的**降级设计**：如果启动时连不上数据库，或某次批量插入抛 SQL 异常，它不会让进程崩溃，而是把 SQL 和这批数据以可读文本打印到控制台（`printToConsole`），生成流程照常推进：

```java
public void batchInsert(String sql, List<Object[]> params) {
    if (!isConnected) {
        printToConsole(sql, params);   // 数据库不可用 → 把"将要执行的 SQL"打出来
        return;
    }
    ...
}
```

Sources: [src/main/datagenerator/business_code/util/DbUtil.java:88-113]()、[DbUtil.java:115-136]()

对教学/演示环境来说，这意味着即使还没有 MySQL，也能先看到程序"会写什么"；对正式跑数仓来说，则要求配置好连接，否则数据只落在控制台而进不了库。

### RandomUtil：仿真感的来源

"像不像真实数据"主要靠 `RandomUtil` 里的词库与取数规则：中文姓名按"张/王/李… × 伟/芳/娜…"拼接，昵称用"快乐/阳光… × 小天使/小星星…"，手机号按 1 开头加 10 位数字，价格保留两位小数，日期往过去随机偏移形成历史分布。分类、品牌、地区也都有固定的候选数组做轮询或随机选取。

Sources: [src/main/datagenerator/business_code/util/RandomUtil.java:15-24]()、[RandomUtil.java:26-78]()、[RandomUtil.java:100-106]()

## 用户日志生成器：把嵌套 JSON 发往 Kafka

### 一条日志的结构

`UserLogGenerator.generateLog()` 每次组装一条完整的日志对象，模型用 `UserLog` 承载，包含 `common`（公共信息：地区/品牌/渠道/设备 ID/用户 ID/OS/版本）、`start`（启动信息：入口/加载时长/广告）、`page`（页面信息：停留时长/页面 ID/来源类型）、`actions`（行为数组）、`displays`（曝光数组）、`err`（错误，多数为空）、`ts`（毫秒时间戳）：

```java
public static UserLog generateLog() {
    UserLog log = new UserLog();
    log.setCommon(generateCommon());
    log.setStart(generateStart());
    log.setPage(generatePage());
    log.setActions(generateActions());
    log.setDisplays(generateDisplays());
    log.setErr(generateError());
    log.setTs(Instant.now().toEpochMilli());
    return log;
}
```

Sources: [src/main/datagenerator/userlog_code/model/UserLog.java:9-15]()；[src/main/datagenerator/userlog_code/generator/UserLogGenerator.java:21-31]()；[src/main/datagenerator/userlog_code/model/Common.java:6-15]()

数组字段的数量也做成随机的：`actions` 是 0–4 条随机行为，`displays` 是 0–9 条曝光，每条含 `item_id`、`item_type`、曝光位等。`err` 只在约 5% 的日志里出现（`random.nextInt(100) < 95` 时为空），模拟偶发错误。

Sources: [src/main/datagenerator/userlog_code/generator/UserLogGenerator.java:68-84]()、[UserLogGenerator.java:86-103]()、[UserLogGenerator.java:105-113]()

### 嵌套数组怎么变成"真 JSON"

这里有个值得一提的实现细节：模型里 `actions`、`displays` 字段声明成 `String`，Jackson 序列化时会把它们当普通字符串，结果是"JSON 里套着被转义的一长串文本"。所以 `UserLogApplication` 发送前用 `processNestedJson` 补救：先把对象变成 JSON 节点，再把 `actions`/`displays` 的字符串内容重新解析成真正的数组节点后写回，最后整体序列化输出：

```java
// 先将对象转换为 JSON 节点
JsonNode rootNode = objectMapper.valueToTree(log);
// 把 actions 字符串解析回 JSON 数组，覆盖原字段
JsonNode actionsNode = objectMapper.readTree(actionsStr);
((ObjectNode) rootNode).set("actions", actionsNode);
```

Sources: [src/main/datagenerator/userlog_code/UserLogApplication.java:105-131]()；[src/main/datagenerator/userlog_code/generator/UserLogGenerator.java:79-83]()、[UserLogGenerator.java:98-102]()

经过这一步，下游拿到的就是 `actions`、`displays` 为原生数组的 JSON，而不是带反斜杠转义的字符串。

### 发送节奏与降级

`UserLogApplication` 启动时先按配置建 `KafkaProducer`（`StringSerializer`、`acks=all`、`retries=3`），并调用 `partitionsFor(topic)` 探活；连不上 Kafka 就标记为本地模式。主循环每 `interval` 毫秒生成一条日志：Kafka 可用则 `send(...).get()` **同步等待发送结果**（保证消息落盘），失败则把该条日志打印到控制台；Kafka 不可用则一直走本地打印。

Sources: [src/main/datagenerator/userlog_code/UserLogApplication.java:45-64]()、[UserLogApplication.java:66-92]()

也就是说，默认配置下它的吞吐并不高（一条一条来、每条约隔 1 秒，同步确认），`userlog_code/README.md` 里"批量发送、异步提高吞吐"的描述和当前代码的实际行为并不一致——代码里是逐条同步发送，以本页代码为准。演示环境可以靠调小 `generator.interval` 来提高数据密度。

Sources: [src/main/datagenerator/userlog_code/UserLogApplication.java:71-91]()；[src/main/datagenerator/userlog_code/README.md:108-122]()

## 小结

整套仿真数据的源头是 `src/main/datagenerator` 下两个独立的 Spring Boot 进程：`BusinessApplication` 用"先打基础数据、再按固定间隔无限循环"的方式，把电商交易链路上的分类、商品、用户、营销、订单、购物车、仓储等几十张表的数据批量写进 MySQL 的 `gmall` 库——主键通过 `MAX(id)` 续接、跨表靠随机数区间近似关联、数据库不可用时降级为控制台打印 SQL；`UserLogApplication` 则把 `common/start/page/actions/displays/err/ts` 组成的嵌套 JSON 日志逐条同步发往 Kafka 主题，并在发送前把嵌套数组还原成真正的 JSON 数组，Kafka 不可用时同样降级为本地打印。MySQL 里的业务表随后由 Flink CDC 同步进 ODS，Kafka 里的日志流由 Flink Kafka Connector 接入 ODS，两条数据线在此汇合，共同支撑后续 DIM/DWD/DWS/ADS 的数仓分层。理解这两个生成器，就能明白湖仓每一条"看起来像真实电商"的明细数据从何而来，也清楚它们的节奏、关联强度与降级边界。

Sources: [README.md:1-11]()；[src/main/datagenerator/business_code/BusinessApplication.java:57-78]()；[src/main/datagenerator/userlog_code/UserLogApplication.java:44-92]()；[src/main/java/com/fy/warehouse/ODS/ods_log_inc.java:24-45]()；[src/main/java/com/fy/warehouse/ODS/ods_sku_info_full.java:21-47]()
