连接注册表:11 种外部数据源的统一接入
在企业数字化转型中,一个中型企业通常拥有 15-30 个独立数据源:MySQL 生产库、PostgreSQL 分析库、MongoDB 日志库、S3 对象存储、Kafka 消息队列、第三方 REST API……
连接注册表:11 种外部数据源的统一接入
“系列:S4 本体建模 · 第 13 篇 | 难度:中级 | 阅读时间:18 分钟
#TL;DR
- ConnectionRegistry 是 coomia-dip 的数据源统一接入层,支持 11 种外部数据源类型(关系型数据库、NoSQL、对象存储、消息队列、API 等),通过声明式配置实现"一次注册、处处可用"。
- 凭证管理与连接测试是生产可靠性的两大基石——平台内建凭证加密存储、定期轮换、连接健康检查三大能力,避免"连接配置散落在各处"的运维噩梦。
- 数据源映射(Source Mapping)将外部数据结构映射为 Ontology ObjectType,让已有系统无需迁移数据即可融入 Ontology 语义体系。
#1. 引言:数据孤岛的统一入口
在企业数字化转型中,一个中型企业通常拥有 15-30 个独立数据源:MySQL 生产库、PostgreSQL 分析库、MongoDB 日志库、S3 对象存储、Kafka 消息队列、第三方 REST API……
传统做法是为每个数据源写一套连接代码、一套认证逻辑、一套健康检查。这导致:
- 连接配置散落在代码仓库、配置文件、环境变量中
- 凭证管理混乱,密码明文出现在配置文件
- 连接失败时缺乏统一的告警和自愈机制
- 新增数据源需要开发新的适配器代码
传统方式: ConnectionRegistry 方式:
App A → MySQL config (硬编码) ┌─────────────────────┐
App B → MySQL config (另一份) │ ConnectionRegistry │
App C → PostgreSQL config │ │
App D → MongoDB config │ 统一注册 │
App E → S3 config │ 统一认证 │
App F → Kafka config │ 统一监控 │
│ 统一映射 │
6 套配置、6 套认证、6 套监控 └─────────────────────┘
│
一套配置、一套认证、一套监控
ConnectionRegistry 就是 coomia-dip 解决这个问题的答案。
#2. ConnectionRegistry 核心概念
#2.1 架构定位
ConnectionRegistry 位于 Control Layer(Control Layer)中,是 SchemaRegistry 的"数据桥梁":
┌─────────────────────────────────────────────────────┐
│ Control Layer (B) │
│ │
│ ┌──────────────┐ ┌──────────────────────┐ │
│ │SchemaRegistry │◄────►│ ConnectionRegistry │ │
│ │ │ │ │ │
│ │ ObjectType │ │ Connection configs │ │
│ │ RelationType │ │ Credential vault │ │
│ │ ActionType │ │ Health monitor │ │
│ │ SourceMapping│ │ Source mapping │ │
│ └──────────────┘ └──────────┬───────────┘ │
│ │ │
└────────────────────────────────────┼────────────────┘
│ gRPC
┌────────────────┼────────────────┐
│ │ │
┌────▼───┐ ┌─────▼────┐ ┌─────▼────┐
│ MySQL │ │PostgreSQL│ │ MongoDB │
│ Oracle │ │ Doris │ │ Redis │
│ SQL Svr │ │ │ │ ES │
└─────────┘ └──────────┘ └──────────┘
#2.2 Connection 数据模型
一个 Connection 由四部分组成:
apiVersion: ontology/v1
kind: Connection
metadata:
name: manufacturing-mysql
namespace: factory-alpha
labels:
environment: production
team: data-engineering
spec:
# 1. 连接类型
type: MYSQL
# 2. 连接参数
config:
host: db-prod-01.internal.company.com
port: 3306
database: manufacturing
charset: utf8mb4
maxPoolSize: 20
connectionTimeout: 5000
# 3. 凭证引用(不直接存储密码)
credentialRef:
name: mysql-prod-credentials
vault: platform-vault
# 4. 健康检查配置
healthCheck:
enabled: true
interval: 60s
timeout: 5s
query: "SELECT 1"
status:
connected: true
lastChecked: "2026-03-24T10:30:00Z"
latency: 12ms
version: "MySQL 8.0.35"
#2.3 Connection 生命周期
┌─────────┐ register ┌──────────┐ test ┌──────────┐
│ PENDING ├─────────────►│ TESTING ├──────────►│ ACTIVE │
└─────────┘ └────┬─────┘ └────┬─────┘
│ fail │
▼ disconnect
┌──────────┐ │
│ FAILED │ ▼
└──────────┘ ┌──────────┐
│ INACTIVE │
└────┬─────┘
│ reconnect
▼
┌──────────┐
│ ACTIVE │
└──────────┘
#3. 11 种支持的数据源类型
#3.1 关系型数据库(4 种)
| 数据源 | 类型标识 | 默认端口 | 驱动 |
|---|---|---|---|
| MySQL | MYSQL | 3306 | mysql-connector-j 8.x |
| PostgreSQL | POSTGRESQL | 5432 | postgresql 42.x |
| Oracle | ORACLE | 1521 | ojdbc11 |
| SQL Server | SQLSERVER | 1433 | mssql-jdbc 12.x |
关系型数据库的连接配置示例:
from ontology_sdk import OntologyClient, ConnectionSpec
client = OntologyClient(base_url="http://control-Layer:8080")
# MySQL 连接
mysql_conn = ConnectionSpec(
name="erp-mysql",
type="MYSQL",
config={
"host": "mysql-prod.internal",
"port": 3306,
"database": "erp_production",
"charset": "utf8mb4",
"maxPoolSize": 20,
"connectionTimeout": 5000,
"ssl": True,
"sslMode": "VERIFY_IDENTITY",
},
credential_ref="erp-mysql-cred",
)
result = client.connection.register(mysql_conn)
print(f"Registered: {result.name}, status: {result.status}")
# PostgreSQL 连接(带 Schema 指定)
pg_conn = ConnectionSpec(
name="analytics-pg",
type="POSTGRESQL",
config={
"host": "pg-analytics.internal",
"port": 5432,
"database": "analytics",
"schema": "public",
"maxPoolSize": 30,
"statementTimeout": 30000,
"ssl": True,
},
credential_ref="pg-analytics-cred",
)
# Oracle 连接(Service Name 模式)
oracle_conn = ConnectionSpec(
name="legacy-oracle",
type="ORACLE",
config={
"host": "oracle-prod.internal",
"port": 1521,
"serviceName": "ORCL",
"maxPoolSize": 10,
"connectionTimeout": 10000,
},
credential_ref="oracle-prod-cred",
)
# SQL Server 连接
sqlserver_conn = ConnectionSpec(
name="hr-sqlserver",
type="SQLSERVER",
config={
"host": "sqlsvr-hr.internal",
"port": 1433,
"database": "HumanResources",
"instanceName": "MSSQLSERVER",
"encrypt": True,
"trustServerCertificate": False,
},
credential_ref="sqlsvr-hr-cred",
)
#3.2 NoSQL 数据库(3 种)
| 数据源 | 类型标识 | 协议 | 特点 |
|---|---|---|---|
| MongoDB | MONGODB | mongodb:// | 文档存储,Schema 自动推断 |
| Redis | REDIS | redis:// | 键值存储,用于缓存映射 |
| Elasticsearch | ELASTICSEARCH | HTTP/HTTPS | 全文搜索,用于搜索索引映射 |
# MongoDB 连接(副本集模式)
mongo_conn = ConnectionSpec(
name="iot-mongodb",
type="MONGODB",
config={
"hosts": [
"mongo-01.internal:27017",
"mongo-02.internal:27017",
"mongo-03.internal:27017",
],
"database": "iot_events",
"replicaSet": "rs0",
"authSource": "admin",
"readPreference": "secondaryPreferred",
"maxPoolSize": 50,
},
credential_ref="mongo-iot-cred",
)
# Elasticsearch 连接
es_conn = ConnectionSpec(
name="log-elasticsearch",
type="ELASTICSEARCH",
config={
"hosts": [
"https://es-01.internal:9200",
"https://es-02.internal:9200",
],
"indexPrefix": "app-logs-",
"numberOfShards": 5,
"numberOfReplicas": 1,
},
credential_ref="es-log-cred",
)
#3.3 对象存储(1 种)
| 数据源 | 类型标识 | 协议 | 文件格式支持 |
|---|---|---|---|
| S3/MinIO | S3 | S3 API | Parquet, CSV, JSON, Avro, ORC |
# S3 / MinIO 连接
s3_conn = ConnectionSpec(
name="data-lake-s3",
type="S3",
config={
"endpoint": "https://s3.amazonaws.com",
"region": "us-east-1",
"bucket": "company-data-lake",
"pathPrefix": "raw/",
"fileFormat": "PARQUET",
"compressionCodec": "SNAPPY",
},
credential_ref="s3-datalake-cred",
)
#3.4 消息队列(1 种)
| 数据源 | 类型标识 | 协议 | 用途 |
|---|---|---|---|
| Kafka | KAFKA | Kafka Protocol | 流式数据接入,CDC |
# Kafka 连接
kafka_conn = ConnectionSpec(
name="cdc-kafka",
type="KAFKA",
config={
"bootstrapServers": "kafka-01:9092,kafka-02:9092,kafka-03:9092",
"groupId": "coomia-dip-consumer",
"autoOffsetReset": "earliest",
"securityProtocol": "SASL_SSL",
"saslMechanism": "SCRAM-SHA-256",
"schemaRegistryUrl": "http://schema-registry:8081",
},
credential_ref="kafka-cdc-cred",
)
#3.5 外部 API(2 种)
| 数据源 | 类型标识 | 协议 | 认证方式 |
|---|---|---|---|
| REST API | REST_API | HTTP/HTTPS | Bearer Token, API Key, OAuth2 |
| gRPC Service | GRPC_SERVICE | gRPC | mTLS, Token |
# REST API 连接
rest_conn = ConnectionSpec(
name="weather-api",
type="REST_API",
config={
"baseUrl": "https://api.weather.example.com/v2",
"authType": "BEARER_TOKEN",
"timeout": 10000,
"retryPolicy": {
"maxRetries": 3,
"backoffMs": 1000,
},
"rateLimiting": {
"requestsPerSecond": 10,
"burstSize": 20,
},
"headers": {
"Accept": "application/json",
},
},
credential_ref="weather-api-cred",
)
# gRPC Service 连接
grpc_conn = ConnectionSpec(
name="pricing-service",
type="GRPC_SERVICE",
config={
"host": "pricing-service.internal",
"port": 9090,
"useTls": True,
"protoPackage": "com.company.pricing.v1",
"serviceName": "PricingService",
"loadBalancingPolicy": "round_robin",
},
credential_ref="pricing-grpc-cred",
)
#4. 连接测试机制
#4.1 注册时测试
每个数据源类型都有内建的连接测试逻辑:
# 注册并测试连接
result = client.connection.register_and_test(mysql_conn)
print(f"Connection: {result.name}")
print(f"Status: {result.status}") # ACTIVE or FAILED
print(f"Latency: {result.latency_ms}ms") # 首次连接延迟
print(f"Version: {result.server_version}") # 服务器版本
print(f"Features: {result.features}") # 支持的特性列表
# 测试结果详情
if result.status == "FAILED":
print(f"Error: {result.error_message}")
print(f"Error Code: {result.error_code}")
print(f"Suggestion: {result.suggestion}")
各数据源类型的测试方法:
数据源类型 测试方法 检查项
──────────────────────────────────────────────────────────
MYSQL SELECT 1 连通性、认证、权限
POSTGRESQL SELECT 1 连通性、Schema 存在性
ORACLE SELECT 1 FROM DUAL 连通性、Service Name
SQLSERVER SELECT 1 连通性、实例名
MONGODB db.runCommand({ping:1}) 连通性、副本集状态
REDIS PING 连通性、认证
ELASTICSEARCH GET /_cluster/health 集群状态、索引访问
S3 HeadBucket 桶存在性、权限
KAFKA listTopics() Broker 连通性
REST_API GET /health 端点可达性
GRPC_SERVICE grpc.health.v1.Check 服务健康状态
#4.2 运行时健康检查
ConnectionRegistry 持续监控所有活跃连接的健康状态:
# 配置健康检查策略
health_config = HealthCheckConfig(
enabled=True,
interval=60, # 每 60 秒检查一次
timeout=5, # 超时 5 秒判定失败
failure_threshold=3, # 连续 3 次失败标记为 INACTIVE
success_threshold=1, # 1 次成功恢复为 ACTIVE
alert_channels=["slack-ops", "pagerduty"],
)
client.connection.update_health_check("erp-mysql", health_config)
健康检查状态机:
success
┌──────────────┐
│ │
▼ │
┌──────────┐ ┌──────────┐
│ HEALTHY │ │DEGRADED │
└─────┬────┘ └────┬─────┘
│ fail │ fail (threshold)
▼ ▼
┌──────────┐ ┌──────────┐
│ CHECKING │ │UNHEALTHY │
└──────────┘ └────┬─────┘
│ auto-reconnect
▼
┌──────────┐
│RECOVERING│
└──────────┘
#4.3 连接池管理
每个数据源连接都配备连接池,避免频繁创建和销毁连接:
# 查看连接池状态
pool_status = client.connection.get_pool_status("erp-mysql")
print(f"Active connections: {pool_status.active}")
print(f"Idle connections: {pool_status.idle}")
print(f"Waiting requests: {pool_status.waiting}")
print(f"Total created: {pool_status.total_created}")
print(f"Total destroyed: {pool_status.total_destroyed}")
print(f"Avg acquire time: {pool_status.avg_acquire_ms}ms")
#5. 凭证管理
#5.1 凭证存储架构
ConnectionRegistry 从不直接存储密码。所有凭证通过 CredentialVault 统一管理:
┌─────────────────────────────────────────────┐
│ ConnectionRegistry │
│ │
│ Connection A ──credentialRef──┐ │
│ Connection B ──credentialRef──┤ │
│ Connection C ──credentialRef──┤ │
│ │ │
│ ┌───────────▼──────────┐ │
│ │ CredentialVault │ │
│ │ │ │
│ │ AES-256-GCM 加密 │ │
│ │ Master Key: HSM/KMS │ │
│ │ Audit Log: 全记录 │ │
│ └──────────────────────┘ │
└─────────────────────────────────────────────┘
#5.2 凭证类型
from ontology_sdk import CredentialSpec
# 用户名/密码凭证
basic_cred = CredentialSpec(
name="mysql-prod-cred",
type="USERNAME_PASSWORD",
data={
"username": "onto_reader",
"password": "encrypted:vault:xxxxx",
},
rotation_policy={
"enabled": True,
"interval_days": 90,
"notify_before_days": 14,
},
)
# API Key 凭证
apikey_cred = CredentialSpec(
name="weather-api-cred",
type="API_KEY",
data={
"apiKey": "encrypted:vault:xxxxx",
"headerName": "X-API-Key",
},
)
# OAuth2 凭证
oauth_cred = CredentialSpec(
name="salesforce-oauth-cred",
type="OAUTH2",
data={
"clientId": "encrypted:vault:xxxxx",
"clientSecret": "encrypted:vault:xxxxx",
"tokenUrl": "https://login.salesforce.com/services/oauth2/token",
"scope": "api refresh_token",
},
)
# mTLS 凭证
mtls_cred = CredentialSpec(
name="pricing-grpc-cred",
type="MTLS",
data={
"certPath": "/certs/client.pem",
"keyPath": "/certs/client-key.pem",
"caPath": "/certs/ca.pem",
},
)
# AWS IAM Role 凭证
iam_cred = CredentialSpec(
name="s3-datalake-cred",
type="AWS_IAM_ROLE",
data={
"roleArn": "arn:aws:iam::123456789:role/coomia-dip-s3-reader",
"externalId": "coomia-dip-prod",
"sessionDuration": 3600,
},
)
#5.3 凭证轮换
生产环境中,凭证必须定期轮换:
# 手动触发凭证轮换
rotation_result = client.credential.rotate("mysql-prod-cred")
print(f"Old credential expired: {rotation_result.old_expired_at}")
print(f"New credential active: {rotation_result.new_active_at}")
print(f"Affected connections: {rotation_result.affected_connections}")
# 查看轮换历史
history = client.credential.rotation_history("mysql-prod-cred")
for entry in history:
print(f" {entry.rotated_at} by {entry.rotated_by} - {entry.status}")
自动轮换流程:
Day 76/90 Day 90/90
(通知即将过期) (自动轮换)
┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ 生成新凭证│──►│双凭证并存 │──►│验证新凭证 │──►│废弃旧凭证│
│ │ │(过渡期) │ │连接正常 │ │记录审计 │
└──────────┘ └──────────┘ └──────────┘ └──────────┘
#5.4 凭证审计
所有凭证操作都会记录审计日志:
# 查询凭证审计日志
audit_logs = client.credential.audit_log(
credential_name="mysql-prod-cred",
since="2026-03-01T00:00:00Z",
)
for log in audit_logs:
print(f"{log.timestamp} | {log.action} | {log.actor} | {log.result}")
# 输出示例:
# 2026-03-15 10:30:00 | READ | service:data-pipeline | SUCCESS
# 2026-03-15 14:22:00 | ROTATE | user:admin@company | SUCCESS
# 2026-03-16 09:00:00 | READ | service:schema-sync | SUCCESS
#6. 数据源映射(Source Mapping)
#6.1 从外部表到 ObjectType
Source Mapping 是 ConnectionRegistry 最核心的功能之一——将外部数据库的表结构映射为 Ontology 的 ObjectType:
# 自动发现外部数据库的表结构
discovery = client.connection.discover_schema("erp-mysql")
for table in discovery.tables:
print(f"Table: {table.name}")
print(f" Columns: {len(table.columns)}")
print(f" Primary Key: {table.primary_key}")
print(f" Foreign Keys: {table.foreign_keys}")
print(f" Row Count: {table.estimated_row_count}")
print()
#6.2 映射配置
apiVersion: ontology/v1
kind: SourceMapping
metadata:
name: erp-equipment-mapping
spec:
connection: erp-mysql
source:
table: equipment
schema: manufacturing
target:
objectType: Equipment
propertyMappings:
- source: equip_id
target: equipmentId
primaryKey: true
- source: equip_name
target: name
- source: equip_status
target: status
transform: "UPPER(value)"
- source: install_date
target: installDate
transform: "CAST(value AS TIMESTAMP)"
- source: line_id
target: productionLineId
foreignKeyMapping:
relation: BelongsToLine
targetType: ProductionLine
targetProperty: lineId
syncPolicy:
mode: INCREMENTAL
schedule: "*/5 * * * *"
watermarkColumn: updated_at
batchSize: 1000
#6.3 映射模式
coomia-dip 支持三种映射模式:
1. VIRTUAL(虚拟映射)
┌──────────┐ 查询时 ┌──────────┐
│Ontology │────────────────►│外部数据库 │
│ObjectType│ 实时转发 │ 原始表 │
└──────────┘ └──────────┘
特点:零延迟、数据不搬迁、依赖外部可用性
2. REPLICATED(复制映射)
┌──────────┐ 定时同步 ┌──────────┐
│Ontology │◄───────────────│外部数据库 │
│Iceberg │ 增量复制 │ 原始表 │
└──────────┘ └──────────┘
特点:高性能查询、有同步延迟、独立可用
3. FEDERATED(联邦映射)
┌──────────┐ 查询路由 ┌──────────┐
│Ontology │────────────────►│外部数据库 │
│查询引擎 │ 带缓存 │ 原始表 │
└──────────┘ └──────────┘
特点:查询优化、智能缓存、兼顾性能和实时性
# 创建虚拟映射
virtual_mapping = SourceMappingSpec(
name="erp-equipment-virtual",
connection="erp-mysql",
source_table="equipment",
target_object_type="Equipment",
mode="VIRTUAL",
property_mappings=[
PropertyMapping("equip_id", "equipmentId", primary_key=True),
PropertyMapping("equip_name", "name"),
PropertyMapping("equip_status", "status", transform="UPPER(value)"),
],
)
# 创建复制映射
replicated_mapping = SourceMappingSpec(
name="erp-equipment-replicated",
connection="erp-mysql",
source_table="equipment",
target_object_type="Equipment",
mode="REPLICATED",
sync_policy=SyncPolicy(
schedule="*/5 * * * *",
watermark_column="updated_at",
batch_size=1000,
conflict_resolution="LATEST_WINS",
),
property_mappings=[
PropertyMapping("equip_id", "equipmentId", primary_key=True),
PropertyMapping("equip_name", "name"),
],
)
client.connection.create_source_mapping(virtual_mapping)
#6.4 Schema 自动推断
对于 MongoDB 等无 Schema 数据源,ConnectionRegistry 提供 Schema 自动推断:
# 从 MongoDB 集合推断 Schema
inferred_schema = client.connection.infer_schema(
connection="iot-mongodb",
collection="sensor_readings",
sample_size=10000, # 采样 10000 条文档
confidence=0.95, # 95% 置信度
)
print(f"Inferred properties: {len(inferred_schema.properties)}")
for prop in inferred_schema.properties:
print(f" {prop.name}: {prop.type} "
f"(nullable: {prop.nullable}, "
f"coverage: {prop.coverage:.1%})")
# 输出示例:
# Inferred properties: 8
# sensorId: STRING (nullable: False, coverage: 100.0%)
# timestamp: TIMESTAMP (nullable: False, coverage: 100.0%)
# temperature: DOUBLE (nullable: True, coverage: 98.5%)
# humidity: DOUBLE (nullable: True, coverage: 97.2%)
# location: STRUCT (nullable: True, coverage: 85.3%)
#7. 连接注册表的 gRPC API
ConnectionRegistry 通过 gRPC 对外提供服务,以下是核心 Protobuf 定义:
syntax = "proto3";
package onto.control.connection.v1;
service ConnectionRegistryService {
// 连接管理
rpc RegisterConnection(RegisterConnectionRequest)
returns (RegisterConnectionResponse);
rpc TestConnection(TestConnectionRequest)
returns (TestConnectionResponse);
rpc UpdateConnection(UpdateConnectionRequest)
returns (UpdateConnectionResponse);
rpc DeleteConnection(DeleteConnectionRequest)
returns (DeleteConnectionResponse);
rpc ListConnections(ListConnectionsRequest)
returns (ListConnectionsResponse);
rpc GetConnectionStatus(GetConnectionStatusRequest)
returns (GetConnectionStatusResponse);
// Schema 发现
rpc DiscoverSchema(DiscoverSchemaRequest)
returns (DiscoverSchemaResponse);
rpc InferSchema(InferSchemaRequest)
returns (InferSchemaResponse);
// 数据源映射
rpc CreateSourceMapping(CreateSourceMappingRequest)
returns (CreateSourceMappingResponse);
rpc UpdateSourceMapping(UpdateSourceMappingRequest)
returns (UpdateSourceMappingResponse);
rpc DeleteSourceMapping(DeleteSourceMappingRequest)
returns (DeleteSourceMappingResponse);
rpc SyncSourceMapping(SyncSourceMappingRequest)
returns (SyncSourceMappingResponse);
// 凭证管理
rpc StoreCredential(StoreCredentialRequest)
returns (StoreCredentialResponse);
rpc RotateCredential(RotateCredentialRequest)
returns (RotateCredentialResponse);
rpc GetCredentialAuditLog(GetCredentialAuditLogRequest)
returns (GetCredentialAuditLogResponse);
// 健康检查
rpc GetHealthStatus(GetHealthStatusRequest)
returns (stream HealthStatusEvent);
}
message RegisterConnectionRequest {
string name = 1;
string namespace = 2;
ConnectionType type = 3;
map<string, string> config = 4;
string credential_ref = 5;
HealthCheckConfig health_check = 6;
}
enum ConnectionType {
CONNECTION_TYPE_UNSPECIFIED = 0;
MYSQL = 1;
POSTGRESQL = 2;
ORACLE = 3;
SQLSERVER = 4;
MONGODB = 5;
REDIS = 6;
ELASTICSEARCH = 7;
S3 = 8;
KAFKA = 9;
REST_API = 10;
GRPC_SERVICE = 11;
}
#8. 多连接编排与查询路由
#8.1 跨数据源查询
当 Ontology 中的 ObjectType 映射到不同数据源时,平台自动处理跨数据源查询:
# Equipment 来自 MySQL,SensorReading 来自 MongoDB,
# WorkOrder 来自 PostgreSQL
# 但用户无需关心数据在哪里
result = client.ontology.query(
object_type="Equipment",
filter="status == 'RUNNING'",
expand=[
"sensors.latestReading", # MongoDB
"workOrders.openCount", # PostgreSQL
],
)
# 平台自动:
# 1. 从 MySQL 查询 Equipment
# 2. 从 MongoDB 查询 SensorReading
# 3. 从 PostgreSQL 查询 WorkOrder
# 4. 在内存中 JOIN 并返回统一结果
#8.2 查询路由策略
┌──────────────────────────────────────────────────────┐
│ Query Router │
│ │
│ 1. 解析查询涉及的 ObjectType │
│ 2. 查找每个 ObjectType 的 SourceMapping │
│ 3. 根据映射模式选择查询路径 │
│ - VIRTUAL → 直接查外部源 │
│ - REPLICATED → 查本地 Iceberg │
│ - FEDERATED → 检查缓存 → 回源 │
│ 4. 并行执行多数据源查询 │
│ 5. 合并结果并应用 Ontology 语义 │
│ │
│ 优化策略: │
│ - Predicate pushdown(谓词下推) │
│ - Projection pushdown(投影下推) │
│ - Join reordering(连接重排序) │
│ - Result caching(结果缓存) │
└──────────────────────────────────────────────────────┘
#9. 运维与监控
#9.1 连接仪表盘
# 获取所有连接的概览
dashboard = client.connection.dashboard()
print(f"Total connections: {dashboard.total}")
print(f"Active: {dashboard.active}")
print(f"Inactive: {dashboard.inactive}")
print(f"Failed: {dashboard.failed}")
# 按类型分组
for type_group in dashboard.by_type:
print(f"\n{type_group.type}:")
print(f" Count: {type_group.count}")
print(f" Avg latency: {type_group.avg_latency_ms}ms")
print(f" Error rate: {type_group.error_rate:.2%}")
#9.2 连接告警规则
# 配置告警规则
alert_rule = AlertRule(
name="connection-high-latency",
condition="connection.latency_ms > 500",
duration="5m",
severity="WARNING",
channels=["slack-ops"],
message="Connection {connection.name} latency exceeded 500ms",
)
client.connection.create_alert_rule(alert_rule)
# 查看告警历史
alerts = client.connection.get_alerts(
since="2026-03-20T00:00:00Z",
severity="WARNING",
)
for alert in alerts:
print(f"{alert.fired_at} | {alert.connection} | {alert.message}")
#9.3 连接迁移
当需要切换数据源(比如从 MySQL 迁移到 PostgreSQL)时:
# 创建迁移计划
migration = client.connection.create_migration(
source_connection="erp-mysql",
target_connection="erp-postgresql",
strategy="BLUE_GREEN", # 蓝绿切换
validation_queries=[
"SELECT COUNT(*) FROM equipment",
"SELECT MAX(updated_at) FROM equipment",
],
)
# 执行迁移
migration.execute()
# 验证
validation = migration.validate()
print(f"Data consistency: {validation.consistency_check}")
print(f"Row count match: {validation.row_count_match}")
print(f"Latency comparison: {validation.latency_comparison}")
# 切换
migration.switch_over()
#10. 安全最佳实践
#10.1 网络层安全
┌─────────────────────────────────────────────────┐
│ coomia-dip VPC │
│ │
│ ┌──────────────┐ ┌──────────────────┐ │
│ │Control Layer │ │ Connection Pool │ │
│ │ │─────►│ │ │
│ │Connection │ │ TLS 1.3 │ │
│ │Registry │ │ Connection limit │ │
│ └──────────────┘ │ IP whitelist │ │
│ └────────┬─────────┘ │
│ │ │
└─────────────────────────────────┼──────────────┘
│ Encrypted
┌──────▼──────┐
│ External DB │
│ (DMZ/VPN) │
└─────────────┘
#10.2 最小权限原则
# 为连接配置最小权限
permission_config = ConnectionPermission(
connection="erp-mysql",
allowed_operations=["SELECT"], # 只读
allowed_tables=["equipment", "orders"], # 只允许特定表
row_filter="factory_id = 'F001'", # 行级过滤
column_mask={
"employee_ssn": "MASK_LAST_4", # 列级脱敏
"salary": "REDACT",
},
)
client.connection.set_permissions(permission_config)
#10.3 审计合规
所有连接操作都会记录完整的审计轨迹:
审计事件 记录内容
────────────────────────────────────────────
CONNECTION_REGISTERED 谁注册、何时、配置快照
CONNECTION_TESTED 测试结果、延迟
CREDENTIAL_ACCESSED 哪个服务读取了凭证
CREDENTIAL_ROTATED 轮换前后的元数据
SCHEMA_DISCOVERED 发现了哪些表
SOURCE_MAPPING_CREATED 映射配置详情
QUERY_ROUTED 查询路由到了哪个连接
CONNECTION_FAILED 失败原因、影响范围
#11. 实战:从零搭建一个多数据源 Ontology
以下是一个完整的实战案例——将三个数据源统一接入 Ontology:
from ontology_sdk import OntologyClient
client = OntologyClient(base_url="http://control-Layer:8080")
# ============================================
# Step 1: 注册数据源
# ============================================
# ERP 系统(MySQL)
client.connection.register(ConnectionSpec(
name="erp-mysql",
type="MYSQL",
config={"host": "mysql.internal", "port": 3306, "database": "erp"},
credential_ref="erp-cred",
))
# IoT 平台(MongoDB)
client.connection.register(ConnectionSpec(
name="iot-mongo",
type="MONGODB",
config={"hosts": ["mongo.internal:27017"], "database": "iot"},
credential_ref="iot-cred",
))
# 数据湖(S3)
client.connection.register(ConnectionSpec(
name="datalake-s3",
type="S3",
config={"endpoint": "https://s3.internal", "bucket": "datalake"},
credential_ref="s3-cred",
))
# ============================================
# Step 2: 测试所有连接
# ============================================
connections = client.connection.list()
for conn in connections:
result = client.connection.test(conn.name)
print(f"{conn.name}: {result.status} ({result.latency_ms}ms)")
# ============================================
# Step 3: 发现 Schema 并创建映射
# ============================================
# 从 MySQL 发现并映射 Equipment
mysql_schema = client.connection.discover_schema("erp-mysql")
equipment_table = mysql_schema.get_table("equipment")
client.connection.create_source_mapping(SourceMappingSpec(
name="equipment-mapping",
connection="erp-mysql",
source_table="equipment",
target_object_type="Equipment",
mode="REPLICATED",
sync_policy=SyncPolicy(schedule="*/5 * * * *"),
property_mappings=equipment_table.auto_map(),
))
# 从 MongoDB 推断并映射 SensorReading
mongo_schema = client.connection.infer_schema("iot-mongo", "sensor_readings")
client.connection.create_source_mapping(SourceMappingSpec(
name="sensor-mapping",
connection="iot-mongo",
source_table="sensor_readings",
target_object_type="SensorReading",
mode="VIRTUAL",
property_mappings=mongo_schema.auto_map(),
))
# 从 S3 映射历史分析数据
client.connection.create_source_mapping(SourceMappingSpec(
name="history-mapping",
connection="datalake-s3",
source_table="analytics/equipment_history/",
target_object_type="EquipmentHistory",
mode="FEDERATED",
property_mappings=[
PropertyMapping("equip_id", "equipmentId"),
PropertyMapping("event_time", "eventTime"),
PropertyMapping("event_type", "eventType"),
PropertyMapping("details", "details"),
],
))
# ============================================
# Step 4: 验证统一查询
# ============================================
# 现在可以跨三个数据源查询
result = client.ontology.query(
object_type="Equipment",
filter="status == 'RUNNING'",
expand=[
"sensorReadings.latest", # 来自 MongoDB
"history.last30Days", # 来自 S3
],
)
for equipment in result.objects:
print(f"Equipment: {equipment.name}")
print(f" Latest sensor: {equipment.sensorReadings.latest}")
print(f" History events: {len(equipment.history.last30Days)}")
#Key Takeaways
-
ConnectionRegistry 是数据源的"统一户口"。11 种数据源类型覆盖了企业常见的所有数据源场景,通过声明式配置实现注册、测试、监控的一体化管理。
-
凭证管理必须独立于连接配置。CredentialVault 提供加密存储、自动轮换、操作审计三大能力,杜绝了密码明文散落各处的安全隐患。
-
Source Mapping 是 Ontology 与外部世界的桥梁。三种映射模式(VIRTUAL/REPLICATED/FEDERATED)适用于不同场景,让已有系统无需迁移数据即可融入 Ontology 语义体系。
-
跨数据源查询对用户透明。Query Router 自动处理谓词下推、结果合并、缓存优化,用户只需关注 Ontology 语义层面的查询。
#下一篇
S4-14: Ontology 建模最佳实践:6 条黄金法则 —— 我们将总结 Ontology 建模过程中的命名规范、粒度把控、关系方向、接口提取、指标设计和版本管理六大最佳实践。
tags: connection-registry, data-source, credential-management, source-mapping, ontology, grpc, connection-testing, health-check, coomia-dip