# data-sync **Repository Path**: shaker/data-sync ## Basic Information - **Project Name**: data-sync - **Description**: No description available - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-09-20 - **Last Updated**: 2026-09-22 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # data-sync 多源数据库 → 可插接 OLAP(ClickHouse / StarRocks / Doris)同步工具。 Rust 实现。默认仅启用 **MySQL 源 + ClickHouse/StarRocks 目标**(纯 Rust 驱动,可离线编译);**Oracle 源**以 feature 门控(需 OCI/ODPI-C 构建环境)。 --- ## 特性 - **可插接连接器**:源(Oracle / MySQL)、目标(ClickHouse / StarRocks / Doris)。新增数据源/目标只需在 `SourceFactory` / `TargetFactory` 注册。 - **中立逻辑类型 + 两段式类型映射**:源原生类型 → `LogicalType` → 目标物理类型,杜绝原生类型泄漏。 - **双同步策略**:快照(snapshot,幂等覆盖分区)+ 增量(incremental,按水位列续传)。 - **结构自治**:自动建表、结构四态判定(CreatedNew / Compatible / Altered / Incompatible)、`auto_alter` 追加列。 - **字段裁剪(R5)**:include / exclude 白黑名单;主键与水位列禁止排除;抽取一律显式列清单,**绝不 `SELECT *`**。 - **隐私友好**:`exclude_columns` 可在抽取阶段剔除非业务字段(ROWID / AUDIT_USER / AUDIT_TS 等),源端审计字段不进入目标 OLAP。 - **调度**:cron 表达式驱动(5/6 字段自动补秒),并发信号量 + 重叠保护。 - **状态存储**:默认纯 Rust 文件存储(零依赖、可离线),可替换为 SQLite / 目标库元数据表。 --- ## 架构 ``` ┌──────────────┐ introspect / read ┌──────────────────┐ 源库 (Oracle/ │ SourceConn │ ───────────────────▶ │ │ MySQL) │ ector │ LogicalColumn/ │ executor (编排) │ └──────────────┘ RowBatch/ChangeBatch │ + schema diff │ │ │ │ effective_columns ▼ ▼ (字段裁剪) ensure_table write_batch / ┌──────────────┐ /alter_add overwrite_scope / 目标库 │ TargetConn │ ◀──────────────────────────────── apply_changes (CH/SR/ │ ector │ DDL + LogicalValue 字面量 (增量 upsert/delete) Doris) └──────────────┘ ▲ StateStore (水位 + 运行记录) ``` 关键模块(`src/`): | 模块 | 职责 | |------|------| | `config.rs` | YAML 配置结构与启动前校验(源引用、cron、增量必带水位列) | | `connector/mod.rs` | 连接器模块入口 | | `connector/source.rs` | 源接口、工厂、字段裁剪逻辑 | | `connector/mysql.rs` | MySQL 源(sqlx,流式快照 + 水位增量) | | `connector/oracle.rs` | Oracle 源(rust-oracle / ODPI-C,阻塞调用经 spawn_blocking) | | `connector/target.rs` | 目标接口、工厂、字面量格式化 | | `connector/clickhouse.rs` | ClickHouse 目标(reqwest HTTP,MergeTree / ReplacingMergeTree) | | `connector/starrocks.rs` | StarRocks / Doris 目标(sqlx MySQL 协议,主键模型 upsert) | | `schema.rs` | 中立逻辑类型、结构 diff、四态判定(纯函数,可单测) | | `executor.rs` | 单次作业编排(快照/增量流程) | | `scheduler.rs` | cron 调度、`once` / `serve` 入口、并发与重叠控制 | | `state.rs` | 运行状态 / 增量水位存储 | | `types.rs` | `LogicalValue`、批次、水位、运行结果等运行时载体 | | `error.rs` | 分层错误类型 | | `log.rs` | 轻量结构化日志(占位,生产可换 tracing) | --- ## 构建 ```bash # 默认:MySQL 源 + ClickHouse/StarRocks 目标,纯 Rust,可离线编译 cargo build --release # 启用 Oracle 源(需 C 编译器构建 ODPI-C,且抽取机安装 Oracle Instant Client) cargo build --release --features oracle # 单元测试(当前 35 个,全绿) cargo test ``` > Oracle 连接器未启用时,配置中引用 `Oracle` 源会返回明确报错提示,而非静默失败。 --- ## 配置 复制 `config.example.yaml` 为 `config.yaml` 后填入真实 DSN。**`config.yaml` 含凭据,已被 `.gitignore` 排除,禁止入库。** ```yaml target: kind: ClickHouse dsn: "http://ch-host:8123" database: traffic_sync sources: - name: oracle_main kind: Oracle dsn: "user/pass@//oracle-host:1521/ORCL" # user/password@connect_string tables: - name: orders_snap source: oracle_main source_table: TRAFFIC.ORDERS target_table: ods_orders strategy: snapshot cron: "0 2 * * *" exclude_columns: [ROWID, AUDIT_USER, AUDIT_TS] # 剔除非业务字段 auto_alter: true on_incompatible: Block - name: orders_incr source: oracle_main target_table: ods_orders_incr strategy: incremental cron: "*/30 * * * *" watermark_column: UPDATE_TS auto_alter: true on_incompatible: Warn - name: mysql_erp kind: MySql dsn: "mysql://user:pass@mysql-host:3306/erp" tables: - name: erp_join source: mysql_erp target_table: ods_erp_join strategy: snapshot cron: "0 4 * * *" query: "SELECT a.id, a.name, b.cat FROM t1 a JOIN t2 b ON a.id=b.id" # 支持 JOIN auto_alter: false on_incompatible: Block scheduler: timezone: Asia/Shanghai max_concurrent_jobs: 4 overlap_protection: true runtime: chunk_size: 50000 state_store: { kind: file, path: "./state/sync_state.json" } ``` 字段说明: - `source_table`:源端表名(可含 owner,如 `TRAFFIC.ORDERS`),为结构自省来源;与 `query` 至少其一必填(配置 `query` 时被自定义 SQL 覆盖)。 - `strategy`:`snapshot`(每日全量,幂等覆盖)| `incremental`(按水位续传)。 - `watermark_column`:增量作业必填,水位列(缺失则配置校验报错)。 - `exclude_columns` / `include_columns`:列裁剪黑名单/白名单;主键与水位列不得被排除。 - `auto_alter`:源新增列时是否自动 `ALTER ADD`;否则遇到新增列终止并报错。 - `on_incompatible`:目标结构不兼容时 `Block`(阻断)| `Warn`(告警继续)。 - `query`:自定义抽取 SQL(覆盖 `source_table`),支持 JOIN。 --- ## 用法 ```bash data-sync demo # 内置演示(纯逻辑,无需数据库):配置校验/字段裁剪/双方言 DDL/状态判定 data-sync validate # 校验 YAML 配置并列出作业 data-sync ddl [--sample] # 生成目标表 DDL(自省源结构或示例结构) data-sync once # 立即运行单个作业一次(需数据库) data-sync serve # 以 cron 守护模式持续运行(需数据库) ``` cron 采用标准 **5 字段**(分 时 日 月 周);程序自动补全秒字段(首位置 0)以适配 6 字段解析器,也可直接写 6 字段。 --- ## 同步策略与字段裁剪 - **快照(snapshot)**:先按 `snapshot_date` 分区幂等删除旧数据,再全量写入;追加 `snapshot_date` / `snapshot_ts` 元数据列。 - **增量(incremental)**:从上次水位(`state/sync_state.json`)或 `MAX(水位列)` 起拉取 `> 水位` 的变更,按主键 upsert/delete;追加 `op` 标记列(ClickHouse 用 ReplacingMergeTree 去重,StarRocks 用主键模型)。 - **字段裁剪(R5)**:`include` 优先,否则 `exclude`;主键与水位列被排除即报错;裁剪后为空即报错;驱动抽取列清单,绝不 `SELECT *`。 ## 目标表状态四态(R3) | 状态 | 含义 | 处理 | |------|------|------| | `CreatedNew` | 目标表不存在 | 自动建表 | | `Compatible` | 结构兼容(仅可空/目标多列差异) | 直接写入 | | `Altered` | 源新增列 | `auto_alter=true` 自动加列,否则终止 | | `Incompatible` | 类型不匹配 / 主键变化 | 依 `on_incompatible` 阻断或告警 | --- ## 安全与隐私 - **凭据**:仅存在于 YAML 配置(DSN),形如 `user/pass@host` 或 `mysql://user:pass@host`。建议通过环境变量 / 密钥管理注入、运行时挂载,不入库。 - **本仓库**:`config.example.yaml` 仅含占位符(`user`/`pass`/`host`),不含任何真实账号、IP 或业务数据。 - **隐私字段**:`exclude_columns` 在抽取阶段剔除 ROWID、审计字段等非业务列,源端敏感/冗余字段不进入目标 OLAP。 - **状态文件**:增量水位与运行记录存于 `state/sync_state.json`,已被 `.gitignore` 排除。 - **已扫描**:源码无硬编码密钥、无 `.env` / 证书 / 私钥;构建日志(`*.log`)仅为 `cargo test` 输出,无 PII。 --- ## 计划与限制 - Oracle 自定义 `query` 结构推断未实现(建议显式列或复用 `introspect`)。 - 增量目前是「基于水位列的简化 upsert」,非 CDC 日志解析(op 固定为 Insert;删除需上游标记)。 - ClickHouse 无主键 upsert,依赖 `ReplacingMergeTree` + `snapshot_ts` 版本去重,查询时需 `FINAL` 或取最新版本。 - 状态存储默认文件实现;生产建议切换为 SQLite / 目标库元数据表以支持多实例。 --- ## 目录结构 ``` data-sync/ ├── Cargo.toml ├── config.example.yaml # 配置模板(占位符,可入库) ├── .gitignore ├── src/ │ ├── main.rs # CLI 入口(demo/validate/ddl/once/serve) │ ├── lib.rs # 库入口 │ ├── config.rs │ ├── connector/ │ │ ├── source.rs # 源接口/工厂/字段裁剪 │ │ ├── target.rs # 目标接口/工厂/字面量 │ │ ├── mysql.rs │ │ ├── oracle.rs │ │ ├── clickhouse.rs │ │ └── starrocks.rs │ ├── schema.rs │ ├── executor.rs │ ├── scheduler.rs │ ├── state.rs │ ├── types.rs │ ├── error.rs │ └── log.rs └── target/ # 构建产物(gitignored) ``` --- ## License 内部工具,未开放源码授权。