返回博客

数据接入指南

数据接入是使用 coomia-dip 平台的第一步。本文介绍如何将关系型数据库、CSV/Excel 文件、API、消息队列和对象存储中的数据导入平台,转化为 Ontology 对象。涵盖连接器配置、Schema 映射、数据验证和常见问题排查。

Coomia发布于 2026年1月22日10 分钟阅读
分享本文Twitter / X

系列:S12 开发者教程 · 第 12 篇 | 难度:入门 | 阅读时间:15 分钟

数据接入指南

#TL;DR

数据接入是使用 coomia-dip 平台的第一步。本文介绍如何将关系型数据库、CSV/Excel 文件、API、消息队列和对象存储中的数据导入平台,转化为 Ontology 对象。涵盖连接器配置、Schema 映射、数据验证和常见问题排查。

#1. 数据接入概述

#1.1 支持的数据源

类别数据源连接方式
关系型数据库MySQL, PostgreSQL, Oracle, SQL ServerJDBC
文件CSV, Excel, JSON, Parquet文件上传 / S3
APIREST API, GraphQLHTTP Connector
消息队列Kafka, RabbitMQ, PulsarStream Connector
对象存储MinIO, S3, OSSS3 Protocol
NoSQLMongoDB, Redis, ElasticsearchNative Driver

#1.2 接入流程

Code
① 注册数据源连接 (Connection)
    ↓
② 创建 Schema 映射 (Mapping)
    ↓
③ 配置数据验证 (Validation)
    ↓
④ 执行数据同步 (Sync)
    ↓
⑤ 验证 Ontology 对象

#2. 关系型数据库接入

#2.1 注册数据库连接

Python
from ontology_sdk import OntoPlatform
from ontology_sdk.connectors import JdbcConnector

platform = OntoPlatform(
    control_plane_url="localhost:50051",
    data_plane_url="localhost:50052"
)

# 注册 MySQL 连接
mysql_conn = platform.connectors.register(
    JdbcConnector(
        name="hr_mysql",
        display_name="HR 系统 MySQL",
        driver="mysql",
        host="hr-db.internal",
        port=3306,
        database="hr_system",
        username="${HR_DB_USER}",
        password="${HR_DB_PASS}",
        connection_pool={
            "min_size": 5,
            "max_size": 20,
            "timeout": 30,
        }
    )
)

# 测试连接
test = platform.connectors.test("hr_mysql")
print(f"连接测试: {'成功' if test.success else '失败'}")
print(f"延迟: {test.latency_ms}ms")
print(f"版本: {test.server_version}")

#2.2 Schema 发现

Python
# 自动发现数据库表结构
discovery = platform.connectors.discover("hr_mysql")

print("=== 发现的表 ===")
for table in discovery.tables:
    print(f"\n表: {table.name} ({table.row_count} 行)")
    for col in table.columns:
        print(f"  {col.name}: {col.type} "
              f"{'NOT NULL' if col.not_null else 'NULLABLE'} "
              f"{'PK' if col.is_primary_key else ''}")

# 建议的 Object Type 映射
suggestions = platform.connectors.suggest_mapping("hr_mysql")
for s in suggestions:
    print(f"\n建议: {s.table_name} -> {s.suggested_object_type}")
    for prop in s.property_mappings:
        print(f"  {prop.column} ({prop.column_type}) -> "
              f"{prop.suggested_property} ({prop.suggested_type})")

#2.3 配置映射并导入

Python
from ontology_sdk.connectors import ImportConfig, FieldMapping

# 配置导入
config = ImportConfig(
    connector="hr_mysql",
    source_table="employees",
    target_object_type="Employee",
    upsert_key="employee_id",
    mappings=[
        FieldMapping("employee_id", "external_id"),
        FieldMapping("full_name", "name"),
        FieldMapping("email_address", "email"),
        FieldMapping("dept_code", "department",
                     lookup={"object_type": "Department", "match_field": "code", "return": "rid"}),
        FieldMapping("base_salary", "salary", transform="round(value, 2)"),
        FieldMapping("hire_date", "hire_date", transform="parse_date(value, 'yyyy-MM-dd')"),
        FieldMapping("job_level", "level",
                     transform={
                         "J1": "junior", "J2": "junior",
                         "M1": "mid", "M2": "mid",
                         "S1": "senior", "S2": "senior",
                     }),
    ],
    filter="status = 'active'",
    batch_size=500,
)

# 执行导入
result = platform.connectors.import_data(config)
print(f"导入完成: {result.total_rows} 行")
print(f"  成功: {result.succeeded}")
print(f"  跳过: {result.skipped}")
print(f"  失败: {result.failed}")
print(f"  耗时: {result.duration_seconds}s")

#3. 文件导入

#3.1 CSV 文件导入

Python
from ontology_sdk.connectors import FileImporter

importer = FileImporter(platform)

# 单文件导入
result = importer.import_csv(
    file_path="/data/customers.csv",
    target_object_type="Customer",
    upsert_key="customer_id",
    encoding="utf-8",
    delimiter=",",
    header_row=1,
    mappings={
        "customer_id": "external_id",
        "company_name": "name",
        "industry": "industry",
        "country": "region",
        "annual_revenue": ("annual_revenue", float),
    },
    skip_empty_rows=True,
    on_error="skip",  # skip / fail / collect
)

print(f"CSV 导入: {result.succeeded}/{result.total_rows} 成功")

# 预览(不实际导入)
preview = importer.preview_csv("/data/customers.csv", limit=5)
for row in preview.rows:
    print(row)
print(f"总行数: {preview.total_rows}")
print(f"列: {preview.columns}")

#3.2 Excel 文件导入

Python
result = importer.import_excel(
    file_path="/data/products.xlsx",
    sheet_name="Sheet1",
    target_object_type="Product",
    upsert_key="sku",
    header_row=1,
    data_start_row=2,
    mappings={
        "A": ("sku", str),           # A 列 -> sku
        "B": ("name", str),          # B 列 -> name
        "C": ("category", str),
        "D": ("price", float),
        "E": ("stock_quantity", int),
    },
)

#3.3 JSON 文件导入

Python
result = importer.import_json(
    file_path="/data/orders.json",
    target_object_type="Order",
    upsert_key="order_id",
    json_path="$.data.orders[*]",  # JSONPath 表达式定位数据数组
    mappings={
        "id": "order_id",
        "customer.id": "customer_external_id",
        "total": ("total_amount", float),
        "status": "status",
        "created_at": ("created_at", "parse_datetime"),
    },
)

#3.4 批量文件导入

Python
# 从目录批量导入 CSV
result = importer.import_directory(
    directory="/data/daily_exports/",
    file_pattern="employees_*.csv",
    target_object_type="Employee",
    upsert_key="employee_id",
    mappings=standard_employee_mapping,
    parallel=4,  # 4 个文件并行处理
)

print(f"处理 {result.file_count} 个文件")
print(f"总记录: {result.total_rows}")
print(f"成功: {result.succeeded}")

#4. API 数据接入

#4.1 REST API 连接器

Python
from ontology_sdk.connectors import ApiConnector, ApiImportConfig

# 注册 API 连接
api_conn = platform.connectors.register(
    ApiConnector(
        name="crm_api",
        display_name="CRM API",
        base_url="https://crm-api.internal/v2",
        auth={
            "type": "bearer",
            "token": "${CRM_API_TOKEN}",
        },
        headers={
            "Content-Type": "application/json",
            "X-API-Version": "2.0",
        },
        rate_limit={
            "requests_per_second": 10,
            "burst": 20,
        },
        timeout_seconds=30,
        retry={
            "max_retries": 3,
            "backoff": "exponential",
        },
    )
)

# 配置 API 数据导入
config = ApiImportConfig(
    connector="crm_api",
    endpoint="/customers",
    method="GET",
    pagination={
        "type": "cursor",
        "cursor_param": "cursor",
        "cursor_path": "$.meta.next_cursor",
        "data_path": "$.data",
        "page_size": 100,
    },
    target_object_type="Customer",
    upsert_key="external_id",
    mappings={
        "id": "external_id",
        "attributes.name": "name",
        "attributes.industry": "industry",
        "attributes.region": "region",
        "attributes.revenue": ("annual_revenue", float),
        "attributes.tier": "tier",
    },
)

result = platform.connectors.import_data(config)
print(f"API 导入: {result.succeeded} 条记录")

#4.2 分页策略

Python
# Offset 分页
pagination_offset = {
    "type": "offset",
    "limit_param": "limit",
    "offset_param": "offset",
    "data_path": "$.results",
    "total_path": "$.total",
    "page_size": 100,
}

# Cursor 分页
pagination_cursor = {
    "type": "cursor",
    "cursor_param": "after",
    "cursor_path": "$.pagination.next_cursor",
    "has_more_path": "$.pagination.has_more",
    "data_path": "$.data",
}

# Page Number 分页
pagination_page = {
    "type": "page_number",
    "page_param": "page",
    "size_param": "per_page",
    "data_path": "$.items",
    "total_pages_path": "$.total_pages",
    "page_size": 50,
}

#5. 消息队列接入

#5.1 Kafka 实时接入

Python
from ontology_sdk.connectors import KafkaConnector

kafka_conn = platform.connectors.register(
    KafkaConnector(
        name="order_kafka",
        display_name="订单事件 Kafka",
        bootstrap_servers="kafka-1:9092,kafka-2:9092",
        topics=["order-created", "order-updated", "order-cancelled"],
        group_id="coomia-dip-order-consumer",
        auto_offset_reset="latest",
        security={
            "protocol": "SASL_SSL",
            "mechanism": "PLAIN",
            "username": "${KAFKA_USER}",
            "password": "${KAFKA_PASS}",
        },
    )
)

# 配置实时同步
platform.connectors.start_stream(
    connector="order_kafka",
    target_object_type="Order",
    upsert_key="order_id",
    event_mapping={
        "order-created": {
            "action": "create",
            "mappings": {
                "payload.id": "order_id",
                "payload.customer_id": "customer_external_id",
                "payload.total": ("total_amount", float),
                "payload.status": "status",
            }
        },
        "order-updated": {
            "action": "update",
            "mappings": {
                "payload.id": "order_id",
                "payload.status": "status",
                "payload.updated_at": "updated_at",
            }
        },
        "order-cancelled": {
            "action": "update",
            "key_field": "payload.id",
            "mappings": {
                "payload.id": "order_id",
                "status": {"value": "cancelled"},
                "payload.cancelled_at": "cancelled_at",
            }
        },
    },
    error_handling={
        "dead_letter_topic": "coomia-dip-dlq",
        "max_retries": 3,
    },
)

#6. 对象存储接入

#6.1 MinIO / S3 接入

Python
from ontology_sdk.connectors import S3Connector

s3_conn = platform.connectors.register(
    S3Connector(
        name="data_lake_minio",
        display_name="数据湖 MinIO",
        endpoint="http://minio:9000",
        access_key="${MINIO_ACCESS_KEY}",
        secret_key="${MINIO_SECRET_KEY}",
        region="us-east-1",
        bucket="data-lake",
    )
)

# 导入 Parquet 文件
result = platform.connectors.import_s3(
    connector="data_lake_minio",
    path="bronze/customers/2025/03/*.parquet",
    target_object_type="Customer",
    upsert_key="customer_id",
    format="parquet",
    mappings={
        "customer_id": "external_id",
        "name": "name",
        "industry": "industry",
    },
)

# 监听新文件自动导入
platform.connectors.watch_s3(
    connector="data_lake_minio",
    path_prefix="bronze/daily/",
    file_pattern="*.csv",
    target_object_type="DailyReport",
    upsert_key="report_id",
    poll_interval=60,  # 每 60 秒检查
    mappings=daily_report_mapping,
)

#7. 数据验证

#7.1 验证规则

Python
from ontology_sdk.connectors import ValidationRule

# 配置验证规则
validation = [
    ValidationRule.not_empty("name", on_fail="reject"),
    ValidationRule.email_format("email", on_fail="reject"),
    ValidationRule.range("salary", min=0, max=10000000, on_fail="flag"),
    ValidationRule.regex("phone", r"^1[3-9]\d{9}$", on_fail="flag"),
    ValidationRule.enum("status", ["active", "inactive", "suspended"], on_fail="reject"),
    ValidationRule.unique("external_id", on_fail="reject"),
    ValidationRule.date_format("hire_date", format="yyyy-MM-dd", on_fail="reject"),
    ValidationRule.not_future("hire_date", on_fail="flag"),
    ValidationRule.custom(
        "age_check",
        lambda row: 18 <= calculate_age(row["birth_date"]) <= 100,
        on_fail="flag",
        message="年龄不在合理范围"
    ),
]

# 应用到导入配置
config.validation_rules = validation
config.validation_threshold = 0.95  # 验证通过率低于 95% 则终止导入

#7.2 验证报告

Python
# 执行数据质量预检
report = platform.connectors.validate(
    connector="hr_mysql",
    source_table="employees",
    rules=validation,
    sample_size=1000,
)

print(f"=== 数据质量报告 ===")
print(f"样本量: {report.sample_size}")
print(f"通过率: {report.pass_rate:.1%}")
print(f"")
for rule_result in report.results:
    status = "PASS" if rule_result.pass_rate >= 0.99 else "WARN" if rule_result.pass_rate >= 0.95 else "FAIL"
    print(f"  [{status}] {rule_result.rule_name}: {rule_result.pass_rate:.1%}")
    if rule_result.failures:
        for f in rule_result.failures[:3]:
            print(f"    行 {f.row}: {f.value} - {f.reason}")

#8. 连接器管理

#8.1 查看和管理连接器

Python
# 列出所有连接器
connectors = platform.connectors.list()
for conn in connectors:
    print(f"{conn.name} [{conn.type}] - {conn.display_name}")
    print(f"  状态: {conn.status}")
    print(f"  最近同步: {conn.last_sync_at}")

# 健康检查
health = platform.connectors.health_check_all()
for name, status in health.items():
    icon = "OK" if status.healthy else "ERR"
    print(f"  [{icon}] {name}: {status.message} ({status.latency_ms}ms)")

# 查看同步历史
history = platform.connectors.get_sync_history("hr_mysql", limit=10)
for run in history:
    print(f"  {run.started_at} | {run.status} | {run.records_synced} 条 | {run.duration}s")

#9. 常见问题排查

#9.1 问题诊断清单

问题可能原因解决方法
连接超时网络/防火墙检查网络连通性,开放端口
认证失败凭据错误验证用户名密码,检查权限
字符编码错误编码不匹配指定正确编码(UTF-8/GBK)
数据类型转换失败类型不兼容添加 transform 函数
唯一键冲突重复数据检查 upsert_key 配置
外键关系缺失依赖对象未导入先导入依赖的 Object Type
性能慢数据量大/无索引增加 batch_size,添加源表索引

#9.2 调试模式

Python
# 开启调试日志
platform.connectors.set_log_level("hr_mysql", "DEBUG")

# 试运行(不写入)
dry_run = platform.connectors.import_data(config, dry_run=True, limit=10)
print(f"试运行结果:")
for sample in dry_run.sample_output:
    print(f"  {sample}")
print(f"预计导入: {dry_run.estimated_total} 条")

#10. 完整实战:多源数据接入

Python
from ontology_sdk import OntoPlatform
from ontology_sdk.connectors import JdbcConnector, ApiConnector, FileImporter, ImportConfig, FieldMapping

platform = OntoPlatform(
    control_plane_url="localhost:50051",
    data_plane_url="localhost:50052"
)

# === 第一步:创建 Object Types ===
platform.schema.create_object_type("Customer", properties={
    "external_id": {"type": "string", "required": True, "indexed": True},
    "name": {"type": "string", "required": True},
    "industry": {"type": "string"},
    "region": {"type": "string"},
    "annual_revenue": {"type": "decimal"},
    "tier": {"type": "string"},
    "source": {"type": "string"},  # 数据来源标记
})

# === 第二步:MySQL 导入基础客户数据 ===
platform.connectors.register(JdbcConnector(
    name="crm_mysql", driver="mysql",
    host="crm-db", port=3306, database="crm",
    username="${CRM_USER}", password="${CRM_PASS}",
))

mysql_result = platform.connectors.import_data(ImportConfig(
    connector="crm_mysql",
    source_table="customers",
    target_object_type="Customer",
    upsert_key="external_id",
    mappings=[
        FieldMapping("id", "external_id"),
        FieldMapping("company_name", "name"),
        FieldMapping("industry_code", "industry"),
        FieldMapping("region", "region"),
    ],
    extra_properties={"source": "crm_mysql"},
))
print(f"MySQL 导入: {mysql_result.succeeded} 条")

# === 第三步:API 补充收入数据 ===
platform.connectors.register(ApiConnector(
    name="billing_api", base_url="https://billing.internal/api",
    auth={"type": "bearer", "token": "${BILLING_TOKEN}"},
))

api_result = platform.connectors.import_data(ImportConfig(
    connector="billing_api",
    endpoint="/customers/revenue",
    target_object_type="Customer",
    upsert_key="external_id",
    mappings=[
        FieldMapping("customer_id", "external_id"),
        FieldMapping("total_revenue", "annual_revenue", transform="float"),
    ],
    merge_mode="update_only",  # 只更新已存在的记录
))
print(f"API 补充: {api_result.succeeded} 条")

# === 第四步:CSV 补充分层数据 ===
importer = FileImporter(platform)
csv_result = importer.import_csv(
    file_path="/data/customer_tiers.csv",
    target_object_type="Customer",
    upsert_key="external_id",
    mappings={
        "customer_id": "external_id",
        "tier": "tier",
    },
    merge_mode="update_only",
)
print(f"CSV 补充: {csv_result.succeeded} 条")

# === 第五步:验证 ===
total = platform.oql.execute("FIND Customer AGGREGATE COUNT(*) AS total")
print(f"\n最终客户总数: {total[0].total}")

quality = platform.oql.execute("""
    FIND Customer
    AGGREGATE
        COUNT(*) AS total,
        COUNT(CASE WHEN annual_revenue IS NOT NULL THEN 1 END) AS has_revenue,
        COUNT(CASE WHEN tier IS NOT NULL THEN 1 END) AS has_tier
""")
print(f"收入数据完整率: {quality[0].has_revenue / quality[0].total:.1%}")
print(f"分层数据完整率: {quality[0].has_tier / quality[0].total:.1%}")

#Key Takeaways

  1. 多源支持:关系型数据库、文件、API、消息队列、对象存储全覆盖
  2. 自动发现:Schema 自动发现和映射建议,减少手工配置
  3. 数据验证:导入前验证数据质量,设定通过率阈值
  4. 增量/实时:支持水位线增量同步和 Kafka 实时接入
  5. 合并策略:多源数据可以 upsert 合并到同一个 Object Type
  6. 调试友好:试运行、预览、调试日志帮助快速排查问题

#Next Article

下一篇:S12-13 订阅与通知指南 — 学习如何配置事件订阅和多渠道通知。

Tags: 数据接入 ETL 连接器 数据导入 Schema映射 数据验证 coomia-dip