数据接入指南
数据接入是使用 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 Server | JDBC |
| 文件 | CSV, Excel, JSON, Parquet | 文件上传 / S3 |
| API | REST API, GraphQL | HTTP Connector |
| 消息队列 | Kafka, RabbitMQ, Pulsar | Stream Connector |
| 对象存储 | MinIO, S3, OSS | S3 Protocol |
| NoSQL | MongoDB, Redis, Elasticsearch | Native 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
- 多源支持:关系型数据库、文件、API、消息队列、对象存储全覆盖
- 自动发现:Schema 自动发现和映射建议,减少手工配置
- 数据验证:导入前验证数据质量,设定通过率阈值
- 增量/实时:支持水位线增量同步和 Kafka 实时接入
- 合并策略:多源数据可以 upsert 合并到同一个 Object Type
- 调试友好:试运行、预览、调试日志帮助快速排查问题
#Next Article
下一篇:S12-13 订阅与通知指南 — 学习如何配置事件订阅和多渠道通知。
Tags: 数据接入 ETL 连接器 数据导入 Schema映射 数据验证 coomia-dip