返回博客

连接注册表:11 种外部数据源的统一接入

在企业数字化转型中,一个中型企业通常拥有 15-30 个独立数据源:MySQL 生产库、PostgreSQL 分析库、MongoDB 日志库、S3 对象存储、Kafka 消息队列、第三方 REST API……

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

连接注册表: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……

传统做法是为每个数据源写一套连接代码、一套认证逻辑、一套健康检查。这导致:

  • 连接配置散落在代码仓库、配置文件、环境变量中
  • 凭证管理混乱,密码明文出现在配置文件
  • 连接失败时缺乏统一的告警和自愈机制
  • 新增数据源需要开发新的适配器代码
Code
传统方式:                               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 的"数据桥梁":

Code
┌─────────────────────────────────────────────────────┐
│                   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 由四部分组成:

YAML
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 生命周期

Code
┌─────────┐   register   ┌──────────┐   test    ┌──────────┐
│ PENDING  ├─────────────►│ TESTING  ├──────────►│  ACTIVE  │
└─────────┘              └────┬─────┘           └────┬─────┘
                              │ fail                  │
                              ▼                  disconnect
                         ┌──────────┐                 │
                         │  FAILED  │                 ▼
                         └──────────┘           ┌──────────┐
                                                │ INACTIVE │
                                                └────┬─────┘
                                                     │ reconnect
                                                     ▼
                                                ┌──────────┐
                                                │  ACTIVE  │
                                                └──────────┘

#3. 11 种支持的数据源类型

#3.1 关系型数据库(4 种)

数据源类型标识默认端口驱动
MySQLMYSQL3306mysql-connector-j 8.x
PostgreSQLPOSTGRESQL5432postgresql 42.x
OracleORACLE1521ojdbc11
SQL ServerSQLSERVER1433mssql-jdbc 12.x

关系型数据库的连接配置示例:

Python
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}")
Python
# 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 种)

数据源类型标识协议特点
MongoDBMONGODBmongodb://文档存储,Schema 自动推断
RedisREDISredis://键值存储,用于缓存映射
ElasticsearchELASTICSEARCHHTTP/HTTPS全文搜索,用于搜索索引映射
Python
# 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/MinIOS3S3 APIParquet, CSV, JSON, Avro, ORC
Python
# 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 种)

数据源类型标识协议用途
KafkaKAFKAKafka Protocol流式数据接入,CDC
Python
# 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 APIREST_APIHTTP/HTTPSBearer Token, API Key, OAuth2
gRPC ServiceGRPC_SERVICEgRPCmTLS, Token
Python
# 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 注册时测试

每个数据源类型都有内建的连接测试逻辑:

Python
# 注册并测试连接
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}")

各数据源类型的测试方法:

Code
数据源类型          测试方法                    检查项
──────────────────────────────────────────────────────────
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 持续监控所有活跃连接的健康状态:

Python
# 配置健康检查策略
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)

健康检查状态机:

Code
                    success
              ┌──────────────┐
              │              │
              ▼              │
┌──────────┐    ┌──────────┐
│  HEALTHY  │    │DEGRADED  │
└─────┬────┘    └────┬─────┘
      │ fail          │ fail (threshold)
      ▼               ▼
┌──────────┐    ┌──────────┐
│ CHECKING  │    │UNHEALTHY │
└──────────┘    └────┬─────┘
                     │ auto-reconnect
                     ▼
               ┌──────────┐
               │RECOVERING│
               └──────────┘

#4.3 连接池管理

每个数据源连接都配备连接池,避免频繁创建和销毁连接:

Python
# 查看连接池状态
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 统一管理:

Code
┌─────────────────────────────────────────────┐
│              ConnectionRegistry              │
│                                             │
│  Connection A ──credentialRef──┐            │
│  Connection B ──credentialRef──┤            │
│  Connection C ──credentialRef──┤            │
│                                │            │
│                    ┌───────────▼──────────┐ │
│                    │   CredentialVault     │ │
│                    │                      │ │
│                    │  AES-256-GCM 加密    │ │
│                    │  Master Key: HSM/KMS │ │
│                    │  Audit Log: 全记录   │ │
│                    └──────────────────────┘ │
└─────────────────────────────────────────────┘

#5.2 凭证类型

Python
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 凭证轮换

生产环境中,凭证必须定期轮换:

Python
# 手动触发凭证轮换
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}")

自动轮换流程:

Code
Day 76/90                     Day 90/90
(通知即将过期)                 (自动轮换)

┌──────────┐   ┌──────────┐   ┌──────────┐   ┌──────────┐
│ 生成新凭证│──►│双凭证并存 │──►│验证新凭证 │──►│废弃旧凭证│
│           │   │(过渡期) │   │连接正常   │   │记录审计  │
└──────────┘   └──────────┘   └──────────┘   └──────────┘

#5.4 凭证审计

所有凭证操作都会记录审计日志:

Python
# 查询凭证审计日志
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:

Python
# 自动发现外部数据库的表结构
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 映射配置

YAML
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 支持三种映射模式:

Code
1. VIRTUAL(虚拟映射)
   ┌──────────┐     查询时      ┌──────────┐
   │Ontology  │────────────────►│外部数据库 │
   │ObjectType│    实时转发      │  原始表   │
   └──────────┘                 └──────────┘
   特点:零延迟、数据不搬迁、依赖外部可用性

2. REPLICATED(复制映射)
   ┌──────────┐     定时同步     ┌──────────┐
   │Ontology  │◄───────────────│外部数据库 │
   │Iceberg   │    增量复制      │  原始表   │
   └──────────┘                 └──────────┘
   特点:高性能查询、有同步延迟、独立可用

3. FEDERATED(联邦映射)
   ┌──────────┐     查询路由     ┌──────────┐
   │Ontology  │────────────────►│外部数据库 │
   │查询引擎   │    带缓存       │  原始表   │
   └──────────┘                 └──────────┘
   特点:查询优化、智能缓存、兼顾性能和实时性
Python
# 创建虚拟映射
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 自动推断:

Python
# 从 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 定义:

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 映射到不同数据源时,平台自动处理跨数据源查询:

Python
# 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 查询路由策略

Code
┌──────────────────────────────────────────────────────┐
│                   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 连接仪表盘

Python
# 获取所有连接的概览
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 连接告警规则

Python
# 配置告警规则
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)时:

Python
# 创建迁移计划
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 网络层安全

Code
┌─────────────────────────────────────────────────┐
│                  coomia-dip VPC                   │
│                                                 │
│  ┌──────────────┐      ┌──────────────────┐    │
│  │Control Layer │      │ Connection Pool   │    │
│  │              │─────►│                  │    │
│  │Connection    │      │ TLS 1.3          │    │
│  │Registry      │      │ Connection limit │    │
│  └──────────────┘      │ IP whitelist     │    │
│                        └────────┬─────────┘    │
│                                 │              │
└─────────────────────────────────┼──────────────┘
                                  │ Encrypted
                           ┌──────▼──────┐
                           │ External DB  │
                           │ (DMZ/VPN)    │
                           └─────────────┘

#10.2 最小权限原则

Python
# 为连接配置最小权限
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 审计合规

所有连接操作都会记录完整的审计轨迹:

Code
审计事件                     记录内容
────────────────────────────────────────────
CONNECTION_REGISTERED        谁注册、何时、配置快照
CONNECTION_TESTED           测试结果、延迟
CREDENTIAL_ACCESSED         哪个服务读取了凭证
CREDENTIAL_ROTATED          轮换前后的元数据
SCHEMA_DISCOVERED           发现了哪些表
SOURCE_MAPPING_CREATED      映射配置详情
QUERY_ROUTED                查询路由到了哪个连接
CONNECTION_FAILED           失败原因、影响范围

#11. 实战:从零搭建一个多数据源 Ontology

以下是一个完整的实战案例——将三个数据源统一接入 Ontology:

Python
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

  1. ConnectionRegistry 是数据源的"统一户口"。11 种数据源类型覆盖了企业常见的所有数据源场景,通过声明式配置实现注册、测试、监控的一体化管理。

  2. 凭证管理必须独立于连接配置。CredentialVault 提供加密存储、自动轮换、操作审计三大能力,杜绝了密码明文散落各处的安全隐患。

  3. Source Mapping 是 Ontology 与外部世界的桥梁。三种映射模式(VIRTUAL/REPLICATED/FEDERATED)适用于不同场景,让已有系统无需迁移数据即可融入 Ontology 语义体系。

  4. 跨数据源查询对用户透明。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