23 KiB
Data Engine 通用数据同步引擎 - 使用文档
一、概述
本系统是一个配置驱动的通用数据同步引擎,核心思想是:平台管理 + 接口管理。
您只需要通过 API 维护好平台和接口的配置信息(认证方式、请求参数、响应解析、目标表结构),系统就会自动:
- 创建目标表(根据
table_definition自动建表) - 拉取数据(分页请求、多步骤请求、并发处理)
- 写入数据库(批量 upsert)
- 增量同步(通过
filtering按最后修改时间过滤) - 自动调度 → 推荐由 PPGo_Job 外部调度(支持内部定时和外部 HTTP 触发两种模式)
- 补偿重试(失败的同步任务自动重试,退避递增)
- Token 自动刷新(OAuth2 token 过期自动续期,不占用重试次数)
不需要为每个平台写一行业务代码。
二、数据库初始化
2.1 核心表
首次使用需执行 sql/init_core_tables.sql,创建以下 4 张表:
| 表名 | 说明 |
|---|---|
api_datasource_platform |
数据源平台配置(认证信息、限流等) |
api_interface |
接口配置(URL、请求参数、表结构等) |
sync_task_log |
同步任务日志(记录状态,用于补偿) |
sync_tracker |
同步跟踪(记录每个接口的最后同步时间) |
2.2 初始化数据
执行 sql/seed_data.sql 创建腾讯广告平台 + 4 个接口配置。
如需清空重来:
ALTER SEQUENCE api_datasource_platform_id_seq RESTART WITH 1;
ALTER SEQUENCE api_interface_id_seq RESTART WITH 1;
DELETE FROM api_interface;
DELETE FROM api_datasource_platform;
\i sql/seed_data.sql
其他平台的 seed 数据:
sql/seed_data_dingtalk.sql— 钉钉平台(部门、用户、考勤)sql/seed_data_dingtalk_hrm.sql— 钉钉 HRMsql/seed_data_dingtalk_salary.sql— 钉钉智能薪酬sql/seed_data_kuaishou.sql— 快手平台sql/seed_data_report_kuaishou.sql— 快手报表引擎演示数据
三、配置说明
3.1 平台管理 (api_datasource_platform)
| 字段 | 说明 | 示例 |
|---|---|---|
platform_code |
平台编码(唯一) | tencent |
platform_name |
平台名称 | 腾讯广告 |
api_base_url |
API 基础地址 | https://api.e.qq.com/v3.0 |
auth_type |
认证类型 | OAUTH2 / TOKEN / API_KEY / SIGN / APP_SIGNATURE |
token |
access_token 明文 | xxxxx |
client_id / client_secret |
OAuth2 凭证 | xxxxx |
auth_config |
自定义认证配置 (JSONB) | 详见 3.1.1 |
rate_limit_per_minute |
每分钟请求限制 | 60 |
request_timeout_ms |
请求超时(毫秒) | 30000 |
注意: 认证凭据现在统一存储在数据库中(
api_datasource_platform表),不再放在config.yml中。config.yml中tencent.oauth.*仅作参考。
3.1.1 auth_config 字段详解
{
"token_in_query": true, // token 放在 URL 查询参数(如腾讯广告)
"query_key": "access_token", // 查询参数名,默认 "access_token"
"header_name": "Authorization", // token 放请求头时的头名
"header_format": "Bearer {token}",
"refresh_token": "xxxxx", // OAuth2 refresh_token
"extra_query_params": { // 额外查询参数
"timestamp": "{timestamp}", // {timestamp} → 当前 Unix 时间戳
"nonce": "{nonce}" // {nonce} → 随机字符串
},
"app_key": "xxx", // SIGN 认证的应用 Key
"app_secret": "xxx", // SIGN 认证的应用 Secret
"sign_algorithm": "md5", // 签名算法: md5 / md5_upper
"sign_secret": "xxx" // 签名专用密钥(可选,默认用 app_secret)
}
3.1.2 认证类型说明
| 类型 | 说明 | token 传递方式 |
|---|---|---|
OAUTH2 |
OAuth2 授权 | URL query 或 Header,支持自动刷新 |
TOKEN |
静态 Token | Authorization: Bearer <token> |
API_KEY |
API Key | X-API-Key <api_key> |
SIGN |
参数签名 | URL query + sign 参数 |
APP_SIGNATURE |
App-ID + 签名 | Header: app-id + signature |
3.2 接口管理 (api_interface)
| 字段 | 说明 | 示例 |
|---|---|---|
platform_id |
所属平台 ID | 1 |
name / code |
接口名称 / 唯一编码 | 图片素材 / image |
url |
接口地址(相对路径) | /images/get |
method |
请求方法 | GET / POST |
request_config |
请求配置 (JSONB) | 详见 3.2.1 |
response_config |
响应配置 (JSONB) | 详见 3.2.2 |
table_definition |
表结构定义 (JSONB) | 详见 3.3 |
3.2.1 request_config 字段详解
{
"parameters_location": "query", // 参数位置: "query"(URL) / "body"(默认POST)
"page": 1,
"page_size": 100,
"page_param": "page", // 分页参数名(可自定义)
"page_size_param": "page_size",
"pagination_mode": "offset", // 分页模式: 默认页码 / "offset"(偏移量)
"time_field": "last_modified_time", // 增量时间字段名
"time_field_mode": "filtering", // 时间模式: "filtering"(腾讯) / "range"(快手)
"cursor_pagination": true, // 游标分页
"initial_cursor": "", // 初始游标值
"fields": ["field1", "field2"], // 请求字段列表
"body_wrapper_field": "param", // Body 包装字段(快手 API 的 param JSON)
"exclude_from_wrapper": ["method"], // 不包装的顶层字段
"row_inject": ["statisticsMonth"], // 将请求参数注入到响应行中
"prefetch": { ... }, // 预取配置(见下文)
"recursive": { ... }, // 递归配置(见下文)
"max_recursive_depth": 20, // 递归最大深度
"full_sync_start_time": 1700000000000 // 全量同步起始时间(毫秒时间戳)
}
parameters_location 说明:
- 未设置或
"body"→ 参数放在请求体中(POST 请求) "query"→ 参数放在 URL 查询字符串中(GET 请求)- 当
method=GET时,即使不设置也默认走 query
支持 4 种分页模式:
| 模式 | 配置 | 说明 |
|---|---|---|
| 普通分页 | 默认 | page 递增,支持 total_page 判断结束 |
| 偏移量分页 | pagination_mode: "offset" |
offset = (page-1) * pageSize |
| 游标分页 | cursor_pagination: true |
从响应 cursor_field 取游标继续下一页 |
| hasMore 分页 | response_config.has_more_field |
从响应 has_more 字段判断是否还有下一页 |
预取(prefetch):某些接口需要"先拿列表 → 遍历每个元素查数据"(如先拉账户列表,再遍历拉图片)。并发处理,并发数由 sync.concurrency 控制。
"prefetch": {
"url": "/advertiser/get", // 预取接口地址
"method": "GET",
"response_path": "data.list", // 从响应中取值路径
"target_param": "account_id", // 注入主请求的参数名
"value_field": "account_id" // 从预取结果取哪个字段值
}
递归(recursive):用于树形结构接口(如钉钉部门树),先查根级 → 对每个子节点递归查下级。
"recursive": {
"key_field": "dept_id", // 递归键字段
"target_param": "dept_id" // 注入到请求的参数名
}
时间分片(range mode):快手等平台使用 beginTime/endTime 做时间范围查询,全量同步时自动按 3 天分片循环拉取。
3.2.2 response_config 字段详解
系统默认解析的响应格式:
{ "code": 0, "message": "success", "data": { "list": [...], "page_info": { "total_page": N } } }
可通过 response_config 自定义:
{
"list_path": "data.list", // 数据列表路径
"success_field": "code", // 成功标识字段,默认 "code"
"success_value": 0, // 成功值,默认 0
"message_field": "message", // 错误消息字段
"cursor_field": "data.cursor", // 游标字段路径(游标分页用)
"has_more_field": "data.has_more", // hasMore 字段路径
"single_record": false // 是否为单条记录(自动包装为数组)
}
3.3 table_definition 字段详解
{
"table_name": "tencent_image",
"columns": [
{ "name": "image_id", "type": "VARCHAR(100)", "comment": "图片ID" },
{ "name": "account_id", "type": "BIGINT", "comment": "账户ID" }
],
"conflict_keys": ["image_id", "account_id"]
}
自动添加的列(无需声明):
| 列名 | 类型 | 说明 |
|---|---|---|
id |
BIGSERIAL PRIMARY KEY | 自增主键 |
tenant_id |
BIGINT | 租户 ID |
creator |
VARCHAR(64) | 创建人 |
created_at |
TIMESTAMPTZ | 创建时间 |
updater |
VARCHAR(64) | 更新人 |
updated_at |
TIMESTAMPTZ | 更新时间 |
deleted_at |
TIMESTAMPTZ | 软删除 |
raw_data |
JSONB | 原始响应数据(含展平的嵌套字段) |
四、API 接口
服务端口:3013(由 config.yml 中 server.address 配置)
路由规则:系统使用 RouteRegister 自动注册,URL 路径由 struct type name 转换为 kebab-case 生成。
4.1 平台管理 CRUD
基础路径:/datasource/platform/controller
| 方法 | 路径 | 说明 |
|---|---|---|
| POST | /datasource/platform/controller/createDatasourcePlatform |
创建平台 |
| GET | /datasource/platform/controller/listDatasourcePlatforms |
平台列表 |
| GET | /datasource/platform/controller/getDatasourcePlatform |
平台详情 |
| GET | /datasource/platform/controller/getPlatformByCode |
按平台编码查询 |
| PUT | /datasource/platform/controller/updateDatasourcePlatform |
更新平台 |
| PUT | /datasource/platform/controller/updateDatasourcePlatformStatus |
更新平台状态 |
| DELETE | /datasource/platform/controller/deleteDatasourcePlatform |
删除平台 |
| GET | /datasource/platform/controller/getPlatformStatistics |
平台统计 |
| POST | /datasource/platform/controller/testPlatformConnection |
测试平台连接 |
4.2 接口管理 CRUD
基础路径:/api/interface/controller
| 方法 | 路径 | 说明 |
|---|---|---|
| POST | /api/interface/controller/createApiInterface |
创建接口 |
| GET | /api/interface/controller/listApiInterfaces |
接口列表 |
| GET | /api/interface/controller/getApiInterface |
接口详情 |
| PUT | /api/interface/controller/updateApiInterface |
更新接口 |
| PUT | /api/interface/controller/updateApiInterfaceStatus |
更新接口状态 |
| DELETE | /api/interface/controller/deleteApiInterface |
删除接口 |
4.3 同步控制
基础路径:/sync/ctrl
| 方法 | 路径 | 说明 |
|---|---|---|
| POST | /sync/ctrl/triggerAll |
全平台同步(PPGo_Job 主任务) |
| POST | /sync/ctrl/compensate |
补偿扫描(PPGo_Job 辅助任务) |
| POST | /sync/ctrl/trigger |
触发单个接口同步 |
| GET | /sync/ctrl/config |
查询平台配置 |
触发单个接口同步示例:
curl -X POST http://localhost:3013/sync/ctrl/trigger \
-H 'Content-Type: application/json' \
-d '{"platformCode":"tencent","interfaceCode":"image","fullSync":true}'
参数说明:
platformCode:平台编码(必填)interfaceCode:接口编码(必填)fullSync:true=全量拉取,false=增量(默认自动判断)
响应:
{
"success": true,
"tableName": "tencent_image",
"totalRows": 1500,
"insertedRows": 1450,
"duration": "12.3s"
}
全平台同步示例:
curl -X POST http://localhost:3013/sync/ctrl/triggerAll -v
响应:
{
"success": true,
"message": "全平台同步完成",
"results": [
{ "platformCode": "tencent", "interfaceCode": "account_relation", "success": true },
{ "platformCode": "tencent", "interfaceCode": "image", "success": true }
]
}
4.4 报表引擎 API
基础路径:/report
| 方法 | 路径 | 说明 |
|---|---|---|
| POST | /report/extractAll |
全量抽取(PPGo_Job 报表任务) |
| POST | /report/extract |
按天数据抽取 |
| POST | /report/backfill |
批量回填数据 |
| POST | /report/query |
用户选择查询 |
| POST | /report/autoCreateTable |
自动创建统计宽表 |
| POST | /report/initTables |
初始化系统表 |
| GET | /report/businesses |
业务列表 |
| GET/POST/DELETE | /report/business[/save] |
业务 CRUD |
| GET | /report/reports |
报表列表 |
| GET/POST/DELETE | /report/report[/save][/saveWithFields] |
报表 CRUD |
| GET | /report/fields |
报表字段列表 |
| GET/POST/DELETE | /report/field[/save] |
字段 CRUD |
| GET | /report/extractConfigs |
抽取配置列表 |
| GET/POST/DELETE | /report/extractConfig[/save] |
抽取配置 CRUD |
4.5 管理页面
| 路径 | 说明 |
|---|---|
GET /admin |
数据引擎管理后台(平台/接口管理) |
GET /admin/report |
报表引擎管理页面 |
五、配置文件 (config.yml)
server:
address: ":3013" # 服务端口
# 报表引擎配置
report:
extract_enabled: false # 是否启用内部定时抽取(推荐 false,由 PPGo_Job 触发)
# 数据同步配置
sync:
page_size: 100 # 每次分页请求条数
concurrency: 5 # 并发处理数(prefetch 遍历实体时)
retry_count: 3 # 最大重试次数
sync_interval_minutes: 60 # 自动同步间隔(内部调度模式用)
compensation_interval_seconds: 300 # 补偿扫描间隔(内部调度模式用)
auto_sync_enabled: false # 是否启用内部自动同步(推荐 false,由 PPGo_Job 触发)
sync_timeout_minutes: 120 # 单次同步超时(分钟)
default_lookback_days: 89 # 全量同步默认回溯天数
default_tenant_id: 1 # 自动同步使用的租户 ID
六、调度策略(核心变更)
系统支持 两种调度模式,通过配置开关切换:
6.1 推荐模式:PPGo_Job 外部调度
config.yml:
report.extract_enabled: false
sync.auto_sync_enabled: false ← 默认值
内部调度器不启动,由 PPGo_Job 定时调用 HTTP 端点触发。详情见第八章。
6.2 备选模式:内部自动调度
config.yml:
report.extract_enabled: true
sync.auto_sync_enabled: true
- 同步引擎:进程启动后自动循环执行全平台同步,每次完成后等待
sync_interval_minutes再继续 - 报表引擎:每天凌晨 2:00 执行全量抽取(gcron
0 0 2 * * *) - 补偿调度器:独立于
auto_sync_enabled,始终随服务启动(内部模式)
6.3 异常中断恢复
- 同步开始前写
sync_tracker.sync_status='running' - 同步完成后写
'success' - 重启检测到
'running'→ 日志告警 → 重新全量同步
七、Token 自动刷新机制
7.1 工作原理
系统支持 OAuth2 token 过期自动刷新,不需要人工介入:
ApiClient发送请求 → 收到 401 或 token 过期错误- 自动调用平台的
RefreshToken回调函数(如RefreshTencentToken) - 刷新成功后立即重试原请求(不消耗退避重试次数)
- 每个
doRequest最多刷新一次,避免死循环 - 刷新后的 token 同时更新内存和数据库
doRequest
├── execute() → HTTP 401 / token_expired
├── isTokenExpiredError() → true
├── RefreshTencentToken()
│ ├── POST https://api.e.qq.com/oauth/refresh_token
│ ├── 更新 platform.Token / platform.AccessToken
│ └── 更新数据库 api_datasource_platform.token + auth_config
└── continue(用新 token 重试,不消耗重试次数)
7.2 token 过期错误检测
系统通过关键字匹配识别 token 过期错误:
tokenExpiredKeywords := []string{
"TOKEN过期", "token过期", "token_expired", "Token过期",
"result=28", // 快手
"access_token", "token已失效", "token无效",
}
7.3 token 过期时的处理策略
| 场景 | 行为 |
|---|---|
| 自动同步中 | 跳过该平台剩余接口,继续处理下一个平台 |
| 补偿扫描中 | 标记为 manual_review,等待人工处理,不再重试 |
| 手动触发 | 自动刷新后重试 |
7.4 腾讯广告 OAuth2 刷新(当前已实现)
已支持腾讯广告 OAuth2 token 自动刷新,使用 refresh_token 换取新的 access_token:
- 刷新端点:
POST https://api.e.qq.com/oauth/refresh_token - 超时:30 秒
- 互斥锁:防止并发刷新(
refreshTokenMu sync.Mutex)
八、PPGo_Job 定时任务配置
系统推荐由 PPGo_Job 外部调度,需在 PPGo_Job 中配置以下任务。
任务 1:全平台数据同步
| 字段 | 值 |
|---|---|
| 任务名称 | 数据引擎 - 全平台同步 |
| 请求 URL | http://<data-engine-host>:3013/sync/ctrl/triggerAll |
| 请求方法 | POST |
| Header | Content-Type: application/json |
| Body | 空(不需要请求体) |
| Cron | 0 0 * * * ?(每小时整点) |
| 超时 | 7200 秒(2小时) |
| 说明 | 遍历所有 ACTIVE 平台下有 table_definition 的接口,自动判断全量/增量。所有平台的同步都在同一任务中完成 |
执行流程:
triggerAll
└── runAutoSync()
├── 查所有 ACTIVE 平台
├── 对每个平台 → 查所有有 table_definition 的接口
│ ├── 查 sync_tracker → 有记录=增量,无记录=全量
│ ├── SyncByConfig(platformCode, interfaceCode)
│ └── token 过期 → 跳过该平台剩余接口
└── 返回每个接口的执行结果
任务 2:失败任务补偿
| 字段 | 值 |
|---|---|
| 任务名称 | 数据引擎 - 补偿扫描 |
| 请求 URL | http://<data-engine-host>:3013/sync/ctrl/compensate |
| 请求方法 | POST |
| Header | Content-Type: application/json |
| Body | 空 |
| Cron | 0 0/5 * * * ?(每5分钟) |
| 超时 | 600 秒(10分钟) |
| 说明 | 扫描 sync_task_log 中 status=failed 的任务,按退避策略重试 |
补偿退避策略:
第1次重试 → 等待 5分钟
第2次重试 → 等待 15分钟
第3次重试 → 等待 30分钟
第4次重试 → 等待 60分钟
第5次重试 → 等待 120分钟(封顶)
达最大次数 → 标记 manual_review,等待人工处理
任务 3:报表每日抽取(可选)
| 字段 | 值 |
|---|---|
| 任务名称 | 数据引擎 - 每日抽取 |
| 请求 URL | http://<data-engine-host>:3013/report/extractAll |
| 请求方法 | POST |
| Header | Content-Type: application/json |
| Body | 空 |
| Cron | 0 0 2 * * ?(每天凌晨2点) |
| 超时 | 3600 秒(1小时) |
| 说明 | 遍历所有启用的业务+报表,抽取前一天数据到统计宽表 |
九、已配置平台
9.1 腾讯广告
| 项目 | 值 |
|---|---|
| platform_code | tencent |
| API Base | https://api.e.qq.com/v3.0 |
| 认证 | OAUTH2(access_token 在 URL query 中,支持自动刷新) |
4 个接口:
| 接口编码 | 名称 | 方法+路径 | 同步模式 | 增量 |
|---|---|---|---|---|
account_relation |
账户列表 | GET /advertiser/get |
普通分页 | 不支持(全量+去重) |
image |
图片素材 | GET /images/get |
预取(遍历账户) | ✅ last_modified_time |
video |
视频素材 | GET /videos/get |
预取(遍历账户) | ✅ last_modified_time |
audio |
音频素材 | POST /muse_audios/get |
普通分页 POST | 不支持(全量+去重) |
图片/视频同步流程:
SyncByConfig("tencent", "image")
├── ① 预取: 分页拉取 /advertiser/get → 所有 account_id
│ 数据同时存入 tencent_account_relation 表
├── ② 并发遍历账户(config.yml concurrency=5)
│ 每个 account_id → GET /images/get?account_id=xxx&page=1&page_size=100
│ 自动补充 filtering(增量时)
│ 结果 upsert 到 tencent_image
└── ③ 更新 sync_tracker 记录同步时间
9.2 其他平台
| 平台 | seed 文件 | 接口数 |
|---|---|---|
| 钉钉 | seed_data_dingtalk.sql |
~10个 |
| 钉钉 HRM | seed_data_dingtalk_hrm.sql |
~5个 |
| 钉钉智能薪酬 | seed_data_dingtalk_salary.sql |
~3个 |
| 快手 | seed_data_kuaishou.sql |
~5个 |
十、补偿机制
10.1 工作原理
- 同步失败 → 自动写入
sync_task_log(status=failed) - 补偿扫描(由 PPGo_Job 调用
/sync/ctrl/compensate)扫描 failed 记录 - 对未达最大重试次数的任务,调用
SyncByConfig重试 - 重试间隔按退避策略递增:5min → 15min → 30min → 60min → 120min
- 达最大次数 → 标记
manual_review,等待人工介入 - Token 过期错误 → 直接标记
manual_review,不再重试
10.2 配置
sync:
compensation_interval_seconds: 300 # 扫描间隔(内部调度用)
retry_count: 3 # 最大重试次数
十一、快速开始
# 1. 建核心表
PGPASSWORD='xxx' psql -h <host> -p 15432 -U postgres -d engine -f sql/init_core_tables.sql
# 2. 初始化腾讯广告数据
PGPASSWORD='xxx' psql -h <host> -p 15432 -U postgres -d engine -f sql/seed_data.sql
# 3. 配置 PPGo_Job(推荐)或启用内部调度
# 见第六章调度策略 和 第八章 PPGo_Job 配置
# 4. 启动服务
cd D:/WorkSpace/HDWL/data-engine
go run main.go
# 5. 手动测试同步
curl -X POST http://localhost:3013/sync/ctrl/trigger \
-H 'Content-Type: application/json' \
-d '{"platformCode":"tencent","interfaceCode":"image","fullSync":true}'
十二、常见问题
Q: 响应格式不符合 {code: 0, data: {list: [...]}} 怎么办?
通过接口管理的 response_config 自定义 list_path、success_field、success_value。如需新的解析逻辑,修改 dynamic_sync.go 的 parseRespExt 函数。
Q: 如何添加新平台?
- 调用平台管理 API 创建平台配置
- 调用接口管理 API 创建接口(含
table_definition) - 平台状态设为
ACTIVE,接口状态设为active - 系统下次同步时自动建表并拉取数据
Q: prefetch 的响应格式要求?
必须是 JSON,response_path 指向一个数组。如 response_path: "data.list" 从 {"data":{"list":[...]}} 取值。
Q: 如何排查同步失败?
- 查询
sync_task_log表看失败记录和错误信息 GET /sync/ctrl/config?platformCode=tencent查看平台配置- 检查服务日志中
equivalent curl输出,手动复现请求 - PPGo_Job 补偿任务会自动重试,日志会打印重试过程
Q: Token 过期了怎么办?
如果配置了 OAuth2 且 refresh_token 有效,系统会自动刷新。手动强制刷新可用 _update_token.go:
# 先修改 _update_token.go 中 token 和 refresh_token 的值
cd D:/WorkSpace/HDWL/data-engine
go run _update_token.go
Q: 内部调度和 PPGo_Job 调度有什么区别?
| 对比项 | 内部调度 | PPGo_Job 调度(推荐) |
|---|---|---|
| 依赖 | 无外部依赖 | 需要 PPGo_Job 服务 |
| 配置 | auto_sync_enabled: true |
auto_sync_enabled: false |
| 控制 | 进程内循环,重启失效 | 统一调度中心,支持暂停/手动执行 |
| 多服务协调 | 不支持 | 支持 |
| 运维 | 需要登录服务器 | 通过 Web UI 管理所有任务 |