1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
|
"""数据资产注册服务"""
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from typing import List, Optional
from datetime import datetime
app = FastAPI(title="数据资产目录服务")
class DataAsset(BaseModel):
asset_id: str
name: str
description: str
domain: str # 业务域
subject: str # 主题域
asset_type: str # TABLE, API, FILE, STREAM
technical_type: str # MySQL, Hive, Kafka, REST...
location: str # 连接信息或路径
owner: str # 数据负责人
owner_team: str
tags: List[str]
sensitivity_level: str # PUBLIC, INTERNAL, CONFIDENTIAL, RESTRICTED
quality_score: Optional[float] = None
created_at: datetime = datetime.now()
class AssetCatalog:
"""资产目录管理"""
def __init__(self, db, neo4j_driver, es_client):
self.db = db
self.graph = neo4j_driver
self.search = es_client
def register_asset(self, asset: DataAsset) -> str:
"""注册新数据资产"""
# 1. 存储到PostgreSQL
self.db.execute("""
INSERT INTO data_assets
(asset_id, name, description, domain, subject,
asset_type, technical_type, location, owner,
owner_team, tags, sensitivity_level, quality_score)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (asset_id) DO UPDATE SET
description = EXCLUDED.description,
quality_score = EXCLUDED.quality_score,
updated_at = NOW()
""", (
asset.asset_id, asset.name, asset.description,
asset.domain, asset.subject, asset.asset_type,
asset.technical_type, asset.location, asset.owner,
asset.owner_team, asset.tags, asset.sensitivity_level,
asset.quality_score
))
# 2. 同步到Elasticsearch(全文检索)
self.search.index(
index="data_assets",
id=asset.asset_id,
document={
"name": asset.name,
"description": asset.description,
"domain": asset.domain,
"tags": asset.tags,
"owner": asset.owner,
"quality_score": asset.quality_score
}
)
return asset.asset_id
def search_assets(self, query: str, filters: dict = None) -> list:
"""全文搜索数据资产"""
es_query = {
"bool": {
"must": [
{"multi_match": {
"query": query,
"fields": ["name^3", "description", "tags^2"]
}}
]
}
}
if filters:
filter_clauses = []
if filters.get("domain"):
filter_clauses.append(
{"term": {"domain": filters["domain"]}}
)
if filters.get("asset_type"):
filter_clauses.append(
{"term": {"asset_type": filters["asset_type"]}}
)
if filter_clauses:
es_query["bool"]["filter"] = filter_clauses
results = self.search.search(
index="data_assets",
query=es_query,
size=20
)
return [hit["_source"] for hit in results["hits"]["hits"]]
|