返回博客

数据接入:从外部数据源到 Ontology 对象

企业数据源的典型分布:

Coomia发布于 2025年8月15日20 分钟阅读
分享本文Twitter / X

数据接入:从外部数据源到 Ontology 对象

系列:S4 本体建模 · 第 11 篇 | 难度:中级 | 阅读时间:18 分钟

#TL;DR

  • 数据接入(Data Onboarding)是将外部数据源映射到 Ontology 模型的过程——CSV 文件、数据库表、REST API、消息队列中的原始数据,通过映射规则转化为 ObjectType 实例,让数据在 Ontology 中"活"起来。
  • **三种同步模式(全量/增量/实时)**覆盖从历史数据迁移到实时数据流的全部场景,每种模式都有明确的适用条件、性能特征和一致性保证。
  • **数据质量门(Data Quality Gate)**在数据进入 Ontology 之前执行验证——类型检查、必填检查、唯一性检查、引用完整性检查,确保"垃圾数据"不会污染 Ontology。

#1. 为什么数据接入是关键挑战

#1.1 数据源的多样性

Code
企业数据源的典型分布:

┌─────────────────────────────────────────────────┐
│             企业数据源全景                        │
│                                                  │
│  结构化数据:                                    │
│  ├── PostgreSQL(订单、客户、产品)              │
│  ├── MySQL(库存、仓储)                         │
│  ├── Oracle(ERP 系统)                          │
│  ├── SQL Server(财务系统)                      │
│  └── CSV/Excel(人工报表、历史数据)             │
│                                                  │
│  半结构化数据:                                  │
│  ├── REST API(第三方服务)                      │
│  ├── GraphQL API(合作伙伴平台)                 │
│  ├── JSON 文件(配置、日志)                     │
│  └── XML 文件(EDI、政府数据)                   │
│                                                  │
│  流数据:                                        │
│  ├── Kafka(事件流)                             │
│  ├── RabbitMQ(消息队列)                        │
│  ├── WebSocket(实时数据推送)                   │
│  └── IoT 设备数据流                              │
│                                                  │
│  问题:                                          │
│  ├── 每种数据源的格式不同                        │
│  ├── 每种数据源的连接方式不同                    │
│  ├── 数据质量参差不齐                            │
│  └── 需要统一映射到 Ontology 模型                │
└─────────────────────────────────────────────────┘

#1.2 coomia-dip 的数据接入架构

Code
数据接入管道:

外部数据源 → Connector → Mapper → Validator → Loader → Ontology

  Connector:连接数据源,读取原始数据
  Mapper:将原始字段映射到 ObjectType 属性
  Validator:执行数据质量检查
  Loader:将验证通过的数据加载到 Ontology

每个步骤都是可配置的、可监控的、可回滚的

#2. Connector:数据源连接

#2.1 数据库 Connector

Code
数据库连接器配置:

connector:
  connectorId: "conn-pg-orders"
  type: DATABASE
  config:
    dialect: POSTGRESQL
    host: "orders-db.internal"
    port: 5432
    database: "orders"
    username: "${ORDERS_DB_USER}"
    password: "${ORDERS_DB_PASS}"
    schema: "public"
    ssl: true
    connectionPool:
      minSize: 2
      maxSize: 10
      idleTimeout: 300s

  tables:
    - tableName: "orders"
      targetObjectType: "Order"
      syncMode: INCREMENTAL
      incrementalColumn: "updated_at"
      primaryKey: "order_id"

    - tableName: "customers"
      targetObjectType: "Customer"
      syncMode: INCREMENTAL
      incrementalColumn: "modified_at"
      primaryKey: "customer_id"

  schedule:
    type: CRON
    expression: "*/5 * * * *"    # 每 5 分钟
    timezone: "Asia/Shanghai"

#2.2 REST API Connector

Code
REST API 连接器配置:

connector:
  connectorId: "conn-api-weather"
  type: REST_API
  config:
    baseUrl: "https://api.weather.com/v3"
    authentication:
      type: API_KEY
      header: "X-API-Key"
      value: "${WEATHER_API_KEY}"
    rateLimit:
      requestsPerSecond: 10
      retryPolicy:
        maxRetries: 3
        backoffMultiplier: 2

  endpoints:
    - path: "/weather/current"
      method: GET
      parameters:
        city: "{{objectId}}"
      targetObjectType: "CityWeather"
      syncMode: FULL
      responseMapping:
        root: "$.data"

  schedule:
    type: INTERVAL
    interval: 15m

#2.3 Kafka Connector

Code
Kafka 流式连接器配置:

connector:
  connectorId: "conn-kafka-events"
  type: KAFKA
  config:
    bootstrapServers: "kafka-1:9092,kafka-2:9092"
    groupId: "coomia-dip-ingestion"
    topics:
      - topic: "order-events"
        targetObjectType: "Order"
        keyField: "orderId"
        format: JSON
      - topic: "user-activities"
        targetObjectType: "UserActivity"
        keyField: "userId"
        format: AVRO
        schemaRegistry: "http://schema-registry:8081"

    consumer:
      autoOffsetReset: EARLIEST
      maxPollRecords: 500
      sessionTimeout: 30s

    syncMode: REALTIME
    errorHandling:
      deadLetterTopic: "coomia-dip-dlq"
      maxRetries: 5

#2.4 CSV/文件 Connector

Code
文件连接器配置:

connector:
  connectorId: "conn-csv-products"
  type: FILE
  config:
    source:
      type: S3
      bucket: "data-imports"
      prefix: "products/"
      filePattern: "*.csv"
    format:
      type: CSV
      delimiter: ","
      header: true
      encoding: "UTF-8"
      quoteChar: '"'
      escapeChar: '\\'
      nullValues: ["", "NULL", "N/A"]

    targetObjectType: "Product"
    syncMode: FULL
    primaryKey: "product_id"

    postProcess:
      archiveProcessed: true
      archivePath: "products/processed/"
      deleteAfterProcess: false

  schedule:
    type: FILE_WATCHER
    pollInterval: 30s

#3. Mapper:字段映射

#3.1 映射规则定义

Code
字段映射配置:

mapping:
  sourceConnector: "conn-pg-orders"
  sourceTable: "orders"
  targetObjectType: "Order"

  fieldMappings:
    # 直接映射(字段名和类型匹配)
    - source: "order_id"
      target: "orderId"
      type: DIRECT

    # 类型转换映射
    - source: "total_amount"
      target: "totalAmount"
      type: CAST
      castConfig:
        from: NUMERIC
        to: DECIMAL
        precision: 18
        scale: 4

    # 枚举映射
    - source: "status"
      target: "orderStatus"
      type: ENUM_MAP
      enumMapping:
        "0": "PENDING"
        "1": "CONFIRMED"
        "2": "SHIPPED"
        "3": "DELIVERED"
        "9": "CANCELLED"

    # 表达式映射
    - source: null
      target: "displayName"
      type: EXPRESSION
      expression: "CONCAT('Order #', order_id, ' - ', customer_name)"

    # 日期格式转换
    - source: "created_at"
      target: "createdAt"
      type: DATE_FORMAT
      sourceFormat: "yyyy-MM-dd HH:mm:ss"
      targetFormat: "ISO8601"

    # JSON 字段提取
    - source: "metadata"
      target: "shippingAddress"
      type: JSON_PATH
      jsonPath: "$.shipping.address"

    # 关系映射
    - source: "customer_id"
      target: "customer"
      type: RELATION
      relationConfig:
        relationType: "OrderBelongsToCustomer"
        targetObjectType: "Customer"
        targetProperty: "customerId"

    # 默认值
    - source: null
      target: "dataSource"
      type: CONSTANT
      value: "orders-db"

#3.2 多源合并

Code
多个数据源映射到同一个 ObjectType:

场景:Customer 数据分散在 3 个系统中

Source 1: CRM 系统(基本信息)
  customer_id → customerId
  full_name → name
  email → email
  phone → phone
  created_date → createdAt

Source 2: 订单系统(消费信息)
  customer_id → customerId(关联键)
  total_orders → orderCount
  total_spent → totalSpent
  last_order_date → lastOrderAt

Source 3: 客服系统(服务信息)
  cust_id → customerId(关联键)
  satisfaction_score → satisfactionScore
  ticket_count → supportTicketCount
  last_contact → lastContactAt

合并策略:
  mergeStrategy:
    primarySource: "crm-system"     # 基本信息以 CRM 为准
    mergeKey: "customerId"          # 通过 customerId 关联
    conflictResolution:
      name: PREFER_PRIMARY           # 名字以主数据源为准
      email: MOST_RECENT             # 邮箱取最近更新的
      phone: PREFER_PRIMARY          # 电话以主数据源为准
      default: MOST_RECENT           # 其他字段取最近更新的

合并后的 Customer 对象同时包含来自 3 个系统的数据
任何一个系统的数据更新都会反映到统一的 Customer 对象上

#3.3 映射验证

Code
映射规则在注册时自动验证:

验证检查项:
├── 源字段存在性:源表/API 中是否存在该字段
├── 目标属性存在性:ObjectType 中是否定义了该属性
├── 类型兼容性:源类型能否安全转换为目标类型
├── 必填覆盖:ObjectType 中的必填属性是否都有映射
├── 主键映射:是否定义了主键的映射规则
├── 关系引用:关联的目标 ObjectType 是否存在
└── 表达式语法:表达式映射的语法是否正确

验证报告示例:
┌─────────────────────────────────────────────────┐
│  Mapping Validation Report                       │
│  Source: orders (PostgreSQL)                     │
│  Target: Order (ObjectType)                      │
├─────────────────────────────────────────────────┤
│  Field mappings: 12                              │
│  Valid: 11                                       │
│  Warnings: 1                                    │
│  Errors: 0                                      │
│                                                  │
│  Warning:                                        │
│  ├── Field "discount_code" in source has no      │
│  │   mapping. Data will be ignored.              │
│                                                  │
│  Coverage:                                       │
│  ├── Source fields mapped: 11/15 (73%)           │
│  ├── Target properties covered: 12/14 (86%)     │
│  ├── Unmapped target: "priority" (has default)   │
│  │   "tags" (nullable)                           │
│  └── All required properties covered: YES        │
└─────────────────────────────────────────────────┘

#4. Validator:数据质量门

#4.1 质量检查规则

Code
数据质量检查的分层策略:

Layer 1 — 格式检查(快速,逐行):
├── 类型匹配:字符串能否解析为目标类型
├── 长度限制:字符串是否超过最大长度
├── 格式匹配:日期、邮箱、URL 等格式验证
├── 范围检查:数值是否在允许范围内
└── 空值检查:必填字段是否为空

Layer 2 — 语义检查(较慢,批量):
├── 唯一性检查:主键是否重复
├── 引用完整性:关联的对象是否存在
├── 枚举值检查:值是否在枚举列表中
├── 业务规则:自定义验证逻辑
└── 跨字段检查:字段间的逻辑关系

Layer 3 — 统计检查(批量完成后):
├── 异常值检测:值是否偏离历史分布
├── 完整性检查:缺失率是否超过阈值
├── 一致性检查:与其他数据源的对比
└── 趋势检查:数值是否有异常波动

#4.2 质量检查配置

Code
质量检查规则配置:

qualityRules:
  objectType: "Order"
  rules:
    - ruleId: "qr-001"
      name: "订单金额非负"
      field: "totalAmount"
      check: RANGE
      config:
        min: 0
        max: 10000000
      severity: ERROR          # 不通过则拒绝

    - ruleId: "qr-002"
      name: "客户ID引用完整"
      field: "customer"
      check: REFERENTIAL_INTEGRITY
      config:
        targetObjectType: "Customer"
        targetField: "customerId"
      severity: ERROR

    - ruleId: "qr-003"
      name: "邮箱格式"
      field: "customerEmail"
      check: REGEX
      config:
        pattern: "^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\\.[a-zA-Z]{2,}$"
      severity: WARNING        # 不通过则标记警告

    - ruleId: "qr-004"
      name: "订单日期合理"
      field: "orderDate"
      check: EXPRESSION
      config:
        expression: "orderDate <= NOW() AND orderDate >= '2020-01-01'"
      severity: ERROR

    - ruleId: "qr-005"
      name: "金额异常检测"
      field: "totalAmount"
      check: ANOMALY
      config:
        method: Z_SCORE
        threshold: 3.0
        baseline: LAST_30_DAYS
      severity: WARNING

#4.3 质量报告

Code
每次数据接入生成质量报告:

┌──────────────────────────────────────────────────────┐
│  Data Quality Report                                  │
│  Source: orders (PostgreSQL)                          │
│  Batch: batch-20250115-001                           │
│  Records: 10,000                                     │
├──────────────────────────────────────────────────────┤
│                                                       │
│  Summary:                                             │
│  ├── Total records: 10,000                            │
│  ├── Passed: 9,823 (98.23%)                          │
│  ├── Warnings: 145 (1.45%)                           │
│  ├── Rejected: 32 (0.32%)                            │
│  └── Quality Score: 98.23 / 100                      │
│                                                       │
│  Rule Results:                                        │
│  ├── qr-001 (金额非负): PASS 9,998 / FAIL 2          │
│  ├── qr-002 (客户引用): PASS 9,970 / FAIL 30         │
│  ├── qr-003 (邮箱格式): PASS 9,855 / WARN 145       │
│  ├── qr-004 (日期合理): PASS 10,000 / FAIL 0         │
│  └── qr-005 (金额异常): PASS 9,980 / WARN 20        │
│                                                       │
│  Rejected Records (sample):                           │
│  ├── Row 3421: totalAmount = -500 (qr-001)           │
│  ├── Row 5892: customer_id = "CUST-9999" not found   │
│  └── Row 7234: totalAmount = -200 (qr-001)           │
│                                                       │
│  Action:                                              │
│  ├── 9,823 records loaded to Ontology                │
│  ├── 145 records loaded with WARNING flag            │
│  ├── 32 records sent to quarantine                   │
│  └── Quarantine records available for manual review  │
└──────────────────────────────────────────────────────┘

#5. 同步模式

#5.1 全量同步(Full Sync)

Code
适用场景:
├── 首次数据导入
├── 小数据量(< 100 万行)
├── 数据源无增量标识
└── 需要完全一致性

流程:
  1. 读取源表全部数据
  2. 映射 + 验证
  3. 比较 Ontology 中已有数据
  4. 执行 INSERT / UPDATE / DELETE
  5. 记录同步结果

比较策略:
  源有 + Ontology 无 → INSERT
  源有 + Ontology 有 + 值变了 → UPDATE
  源有 + Ontology 有 + 值没变 → SKIP
  源无 + Ontology 有 → DELETE(或标记为 STALE)

全量同步配置:
syncConfig:
  mode: FULL
  batchSize: 1000
  parallelism: 4
  deletePolicy: SOFT_DELETE    # 不物理删除,标记为失效
  conflictResolution: SOURCE_WINS
  timeout: 30m

#5.2 增量同步(Incremental Sync)

Code
适用场景:
├── 数据源有 updated_at / version 字段
├── 中大数据量(100 万 ~ 1 亿行)
├── 需要定期同步(每 5 分钟 ~ 每小时)
└── 可以接受短暂不一致

流程:
  1. 读取上次同步的水位线(Watermark)
  2. 查询 updated_at > watermark 的记录
  3. 映射 + 验证
  4. 执行 INSERT / UPDATE
  5. 更新水位线

水位线管理:
  watermarkStore:
    connectorId: "conn-pg-orders"
    tableName: "orders"
    watermarkColumn: "updated_at"
    currentWatermark: "2025-01-15T10:25:00Z"
    lastSyncRecords: 342
    lastSyncDuration: 12.5s

增量同步的陷阱:
├── 删除操作:增量同步无法检测到源表中的 DELETE
│   → 解决方案:软删除字段(is_deleted)或定期全量比较
├── 回填数据:如果有人修改了 updated_at 之前的记录
│   → 解决方案:定期全量校验(如每天一次)
├── 时钟偏差:源数据库的时钟可能不精确
│   → 解决方案:水位线留 overlap(如回退 30 秒)
└── 大批量更新:一次更新 100 万行,增量查询很慢
    → 解决方案:分批次处理 + 超时切换全量模式

#5.3 实时同步(Real-time Sync)

Code
适用场景:
├── 需要秒级延迟
├── 数据源支持 CDC 或消息队列
├── 关键业务数据
└── 不能接受批量延迟

CDC(Change Data Capture)方式:
  1. 监听数据库的 binlog / WAL
  2. 捕获 INSERT / UPDATE / DELETE 事件
  3. 映射 + 验证
  4. 实时写入 Ontology

  配置示例:
  connector:
    type: CDC
    config:
      database: POSTGRESQL
      host: "orders-db.internal"
      replicationSlot: "coomia-dip_cdc"
      publication: "orders_pub"
      tables: ["orders", "customers", "line_items"]

  优点:
  ├── 延迟低(秒级)
  ├── 捕获所有变更(包括 DELETE)
  ├── 不增加源数据库负载(读 WAL,不查表)
  └── 保证顺序性

消息队列方式:
  1. 应用程序发送事件到 Kafka
  2. coomia-dip 消费事件
  3. 映射 + 验证
  4. 写入 Ontology

  延迟更低(毫秒级)
  但需要应用程序配合发送事件

#6. 数据接入编排

#6.1 依赖顺序

Code
多个 ObjectType 的数据接入有顺序依赖:

场景:电商数据接入

依赖关系:
  Customer → 无依赖(先导入)
  Product → 无依赖(先导入)
  Order → 依赖 Customer(Customer 必须先存在)
  LineItem → 依赖 Order + Product
  Payment → 依赖 Order

接入顺序编排:
  Phase 1(并行):Customer, Product
  Phase 2(并行):Order(等 Phase 1 完成)
  Phase 3(并行):LineItem, Payment(等 Phase 2 完成)

编排配置:
orchestration:
  name: "ecommerce-full-import"
  phases:
    - phase: 1
      parallel: true
      connectors:
        - "conn-crm-customers"
        - "conn-pim-products"
    - phase: 2
      dependsOn: [1]
      connectors:
        - "conn-oms-orders"
    - phase: 3
      dependsOn: [2]
      parallel: true
      connectors:
        - "conn-oms-lineitems"
        - "conn-pay-payments"

  errorPolicy:
    phaseFailure: STOP_ALL    # 某个 Phase 失败则停止后续
    connectorFailure: CONTINUE # 同 Phase 中的一个失败不影响其他

#6.2 回填策略

Code
首次接入大量历史数据的回填策略:

场景:接入 5 年的历史订单数据(5000 万行)

策略 1 —— 时间分片回填:
  将 5 年数据按月分片
  每个分片独立导入
  并行处理多个分片

  timeline:
    sliceBy: MONTH
    startDate: "2020-01-01"
    endDate: "2025-01-15"
    parallelSlices: 4

  效果:
  60 个月 × 83 万行/月
  4 个并行分片
  预计耗时:~2 小时

策略 2 —— 优先级回填:
  先导入最近的数据(用户立即可用)
  后台慢慢回填历史数据

  priority:
    - range: "LAST_30_DAYS"    # 先导入最近 30 天
      parallelism: 8
    - range: "LAST_365_DAYS"   # 再导入过去 1 年
      parallelism: 4
    - range: "ALL"             # 最后导入全部历史
      parallelism: 2

策略 3 —— 渐进式回填:
  白天低负载运行(不影响业务)
  夜间高负载运行

  schedule:
    daytime:   # 09:00-18:00
      parallelism: 2
      batchSize: 500
    nighttime: # 18:00-09:00
      parallelism: 8
      batchSize: 5000

#7. 错误处理与隔离

#7.1 隔离区(Quarantine)

Code
验证失败的记录进入隔离区:

隔离区的作用:
├── 不让坏数据进入 Ontology
├── 保留原始数据供人工审查
├── 记录失败原因
├── 支持修复后重新导入
└── 提供质量趋势分析

隔离区记录:
{
  "quarantineId": "q-20250115-001",
  "connectorId": "conn-pg-orders",
  "batchId": "batch-20250115-001",
  "sourceRecord": {
    "order_id": "ORD-999",
    "customer_id": "CUST-9999",
    "total_amount": -500,
    "status": "1"
  },
  "targetObjectType": "Order",
  "failedRules": [
    {
      "ruleId": "qr-001",
      "ruleName": "订单金额非负",
      "field": "totalAmount",
      "value": -500,
      "reason": "Value -500 is below minimum 0"
    },
    {
      "ruleId": "qr-002",
      "ruleName": "客户引用完整",
      "field": "customer_id",
      "value": "CUST-9999",
      "reason": "Customer CUST-9999 not found in Ontology"
    }
  ],
  "quarantinedAt": "2025-01-15T10:30:00Z",
  "status": "PENDING_REVIEW"
}

隔离区操作:
  人工审查 → 修复数据 → 重新验证 → 导入 Ontology
  或
  人工审查 → 确认为垃圾数据 → 永久丢弃

#7.2 重试机制

Code
暂时性错误的重试策略:

暂时性错误类型:
├── 网络超时
├── 数据库连接断开
├── 消息队列不可用
├── API 限流(429)
└── 目标存储暂时不可用

重试配置:
retryPolicy:
  maxRetries: 5
  backoff:
    type: EXPONENTIAL
    initialDelay: 1s
    maxDelay: 60s
    multiplier: 2
  retryableErrors:
    - CONNECTION_TIMEOUT
    - TEMPORARY_UNAVAILABLE
    - RATE_LIMITED
  nonRetryableErrors:
    - VALIDATION_FAILED
    - MAPPING_ERROR
    - AUTHENTICATION_FAILED

重试示例:
  Attempt 1: 失败(超时)→ 等待 1s
  Attempt 2: 失败(超时)→ 等待 2s
  Attempt 3: 失败(超时)→ 等待 4s
  Attempt 4: 成功 ✓

死信队列(DLQ):
  5 次重试都失败 → 消息进入 DLQ
  DLQ 中的消息不会自动重试
  需要人工排查原因后手动重放

#8. 数据接入监控

#8.1 接入状态仪表板

Code
数据接入监控面板:

┌──────────────────────────────────────────────────────┐
│  Data Onboarding Dashboard                            │
├──────────────────────────────────────────────────────┤
│                                                       │
│  Active Connectors: 12 / 15                          │
│  ├── Healthy: 10                                     │
│  ├── Warning: 2 (lag > 5min)                         │
│  └── Error: 0                                        │
│                                                       │
│  Last 24h Summary:                                   │
│  ├── Records processed: 2,345,678                    │
│  ├── Records loaded: 2,310,234 (98.5%)              │
│  ├── Records quarantined: 35,444 (1.5%)             │
│  ├── Average latency: 3.2s                          │
│  └── Peak throughput: 12,500 records/s              │
│                                                       │
│  Connector Health:                                   │
│  ┌──────────────────┬────────┬─────────┬──────────┐ │
│  │ Connector        │ Status │ Lag     │ Quality  │ │
│  ├──────────────────┼────────┼─────────┼──────────┤ │
│  │ pg-orders        │ ✅     │ 30s    │ 99.2%    │ │
│  │ kafka-events     │ ✅     │ 2s     │ 97.8%    │ │
│  │ api-weather      │ ⚠️    │ 8m     │ 100%     │ │
│  │ csv-products     │ ✅     │ 0s     │ 98.5%    │ │
│  │ cdc-inventory    │ ✅     │ 1s     │ 99.9%    │ │
│  └──────────────────┴────────┴─────────┴──────────┘ │
└──────────────────────────────────────────────────────┘

#8.2 血缘追踪

Code
每条 Ontology 数据都可以追溯到数据来源:

数据血缘(Data Lineage):

查询:这个 Order 的数据来自哪里?

GET /api/v1/objects/Order/ORD-001/lineage

{
  "objectId": "ORD-001",
  "objectType": "Order",
  "sources": [
    {
      "connector": "conn-pg-orders",
      "table": "orders",
      "primaryKey": "ORD-001",
      "lastSyncAt": "2025-01-15T10:30:00Z",
      "syncMode": "INCREMENTAL",
      "batchId": "batch-20250115-042"
    }
  ],
  "propertyLineage": {
    "orderId": {"source": "orders.order_id", "transform": "DIRECT"},
    "totalAmount": {"source": "orders.total_amount", "transform": "CAST(DECIMAL)"},
    "orderStatus": {"source": "orders.status", "transform": "ENUM_MAP(0→PENDING,1→CONFIRMED...)"},
    "customer": {"source": "orders.customer_id", "transform": "RELATION(Customer)"}
  },
  "qualityFlags": {
    "customerEmail": "WARNING: invalid format"
  }
}

用途:
├── 数据质量排查:这个奇怪的值是哪来的?
├── 合规审计:这些客户数据的来源是什么?
├── 影响分析:如果源表结构变了,哪些 Ontology 对象受影响?
└── 调试:为什么这个对象的值和源数据不一致?

#9. 实战:电商平台数据接入

#9.1 完整接入方案

Code
电商平台数据接入全景:

┌─────────────────────────────────────────────────────┐
│                                                      │
│  [CRM MySQL] ──CDC──→ Customer (实时)                │
│  [PIM PostgreSQL] ──增量──→ Product (5分钟)          │
│  [OMS PostgreSQL] ──CDC──→ Order, LineItem (实时)    │
│  [支付网关 API] ──增量──→ Payment (1分钟)            │
│  [Kafka事件流] ──实时──→ UserActivity (实时)         │
│  [物流API] ──增量──→ Shipment (15分钟)               │
│  [CSV报表] ──全量──→ FinancialReport (每天)          │
│  [IoT MQTT] ──实时──→ WarehouseDevice (实时)         │
│                                                      │
│  接入顺序:                                          │
│  Phase 1: Customer, Product(基础主数据)             │
│  Phase 2: Order(依赖 Customer)                     │
│  Phase 3: LineItem, Payment, Shipment(依赖 Order)  │
│  Phase 4: UserActivity, WarehouseDevice(独立流)     │
│  Phase 5: FinancialReport(每日汇总)                 │
│                                                      │
│  数据量预估:                                        │
│  Customer: 500 万,增量 ~1000/天                      │
│  Product: 100 万,增量 ~500/天                        │
│  Order: 5000 万,增量 ~50000/天                       │
│  LineItem: 2 亿,增量 ~200000/天                      │
│  UserActivity: ~1000 万/天(事件流)                  │
└─────────────────────────────────────────────────────┘

#10. 与 Palantir Foundry 的对比

Code
数据接入能力对比:

┌──────────────────┬────────────────────┬────────────────────┐
│ 特性              │ coomia-dip           │ Palantir Foundry   │
├──────────────────┼────────────────────┼────────────────────┤
│ Connector 类型    │ DB/API/Kafka/File  │ 200+ connectors    │
│ 映射方式          │ 声明式 YAML        │ Pipeline Builder   │
│ 质量检查          │ 多层验证           │ Checks framework   │
│ CDC 支持          │ PostgreSQL/MySQL   │ 多数据库           │
│ 实时流            │ Kafka/MQTT         │ 多种流平台         │
│ 数据血缘          │ 属性级别           │ 列级别             │
│ 隔离区            │ ✓                  │ ✓                  │
│ 多源合并          │ ✓                  │ ✓                  │
│ 编排              │ Phase-based        │ Pipeline DAG       │
│ 可视化            │ Dashboard          │ Pipeline 可视化    │
│ 社区 Connector    │ 开发中             │ 丰富               │
└──────────────────┴────────────────────┴────────────────────┘

coomia-dip 的数据接入目标不是复制 Foundry 的 200+ connector,
而是提供一个灵活的框架,让 80% 的常见场景
(关系数据库、REST API、Kafka、CSV)开箱即用。

#Key Takeaways

  1. 数据接入是 Ontology 的"生命线"——再完美的 ObjectType 定义,如果没有可靠的数据接入管道,Ontology 就是一个空壳。Connector + Mapper + Validator + Loader 四步流水线确保数据可靠地从外部流入 Ontology。
  2. 三种同步模式覆盖全场景——全量(首次导入、小数据量)、增量(定期同步、中大数据量)、实时(CDC/Kafka,秒级延迟),根据业务需求和数据源能力选择合适的模式。
  3. 数据质量门是核心防线——格式检查、语义检查、统计检查三层验证确保"垃圾数据"被拦截在 Ontology 之外,隔离区保留原始数据供人工审查和修复。
  4. 多源合并实现了统一视图——一个 Customer 的数据可能来自 CRM、订单系统、客服系统三个数据源,合并策略自动处理冲突,消费者看到的是统一的 Customer 对象。
  5. 数据血缘是信任的基础——每条数据都可以追溯到原始来源、映射规则、质量标记,当数据出现问题时,可以快速定位到是哪个数据源、哪次同步引入的。

#Next Article

下一篇 S4-12 自动供给(Auto-Provisioning) 将讨论 Ontology 的自动化基础设施供给——当新的 ObjectType 注册时,系统如何自动创建存储表、索引、API 端点、权限配置,实现"定义即部署"的零人工介入。

#ontology #data-onboarding #etl #connector #data-quality #cdc #incremental-sync #data-lineage #mapping