数据接入:从外部数据源到 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
- 数据接入是 Ontology 的"生命线"——再完美的 ObjectType 定义,如果没有可靠的数据接入管道,Ontology 就是一个空壳。Connector + Mapper + Validator + Loader 四步流水线确保数据可靠地从外部流入 Ontology。
- 三种同步模式覆盖全场景——全量(首次导入、小数据量)、增量(定期同步、中大数据量)、实时(CDC/Kafka,秒级延迟),根据业务需求和数据源能力选择合适的模式。
- 数据质量门是核心防线——格式检查、语义检查、统计检查三层验证确保"垃圾数据"被拦截在 Ontology 之外,隔离区保留原始数据供人工审查和修复。
- 多源合并实现了统一视图——一个 Customer 的数据可能来自 CRM、订单系统、客服系统三个数据源,合并策略自动处理冲突,消费者看到的是统一的 Customer 对象。
- 数据血缘是信任的基础——每条数据都可以追溯到原始来源、映射规则、质量标记,当数据出现问题时,可以快速定位到是哪个数据源、哪次同步引入的。
#Next Article
下一篇 S4-12 自动供给(Auto-Provisioning) 将讨论 Ontology 的自动化基础设施供给——当新的 ObjectType 注册时,系统如何自动创建存储表、索引、API 端点、权限配置,实现"定义即部署"的零人工介入。
#ontology #data-onboarding #etl #connector #data-quality #cdc #incremental-sync #data-lineage #mapping