name: data-fabric-weaver description: > Data Fabric / 資料編織 / 後設資料驅動 / 智慧資料整合 / 主動後設資料 / 資料虛擬化 / 知識圖譜 / Apache Atlas / 資料目錄 / Data Mesh 過渡 / 增強資料湖 / 自助資料服務 / 資料資產地圖 version: 1.0.0 author: Marvis tags: - data-fabric - metadata-driven - data-integration - knowledge-graph - apache-atlas - data-virtualization - active-metadata - data-catalog
用後設資料的絲線編織資料生態——讓異構資料來源無縫對話,讓資料消費像開啟 App Store 一樣簡單。
┌─────────────────────────────────┐
│ Data Fabric 編織層 │
└────────────┬────────────────────┘
│
┌──────────┬──────────┬──────┴──────┬──────────┬──────────┐
▼ ▼ ▼ ▼ ▼ ▼
主動後設資料 知識圖譜 資料虛擬化 自動化編排 自助服務 治理嵌入
(Active (Knowledge (Data (Auto- (Self- (Governed
Metadata) Graph) Virtualization) Orchestr) Service) by Design)
│ │ │ │ │ │
▼ ▼ ▼ ▼ ▼ ▼
Atlas Neo4j Trino/ Airflow/ Data RBAC+
OpenMetadata NebulaGraph Presto Dagster Portal Policy
| 維度 | 傳統ETL/資料湖 | Data Fabric |
|---|---|---|
| 整合方式 | 手動編寫 ETL 管道,逐源接入 | 後設資料驅動,一次註冊,全域性可用 |
| 資料發現 | 翻閱文件/問開發,T+1 級 | 即時目錄檢索,秒級 |
| 查詢方式 | 每個源寫專屬查詢 | 虛擬化統一 SQL,跨源聯邦查詢 |
| 新增資料來源 | 2-4 周(ETL開發+測試) | 1-2 天(註冊後設資料+策略) |
| 治理粒度 | 表級/庫級粗粒度 | 列級/行級細粒度(主動後設資料) |
| 資料消費 | 提需求 → 排期 → 取數 | 自助資料市場,即查即用 |
# 一鍵部署 Data Fabric 實驗環境
docker-compose up -d
# 註冊新資料來源
python scripts/register_datasource.py --type postgres --host localhost --db sales
# 搜尋資料資產
python scripts/search_assets.py --keyword "客戶交易"
# 查詢血緣關係
python scripts/trace_lineage.py --table dwd_sales_order
# 聯邦查詢演示
python scripts/federated_query.py
7w4.net提供免費和付費技能下載。
主動後設資料不僅記錄"資料在哪裡",更自動感知變化並推送建議。
# scripts/datasource_registry.yaml
# 資料來源註冊清單 - 後設資料自動採集
datasources:
- name: "交易資料庫"
type: mysql
host: "10.0.1.100"
port: 3306
tags: ["P0", "金融", "OLTP"]
- name: "使用者行為日誌"
type: kafka
brokers: ["10.0.2.10:9092", "10.0.2.11:9092"]
topic_pattern: "user_behavior_*"
tags: ["P1", "即時", "使用者"]
- name: "資料湖儲存"
type: s3
endpoint: "minio:9000"
bucket: "data-lake"
format: "iceberg"
tags: ["P2", "湖倉一體"]
- name: "Elasticsearch 搜尋日誌"
type: elasticsearch
host: "10.0.3.20"
port: 9200
index_pattern: "logs-*"
tags: ["P3", "日誌", "全文檢索"]
# 血緣關係自動推導
lineage_inference:
enabled: true
sql_parser: "sqlglot" # 解析所有 SQL 自動推導血緣
automatic_discovery: true # 自動發現上下游依賴
# scripts/active_metadata.py
"""主動後設資料採集引擎"""
from atlasclient.client import Atlas
import sqlglot
import json
class ActiveMetadataCollector:
"""
主動後設資料採集器
功能:自動掃描資料來源、解析SQL推導血緣、生成智慧標籤
"""
def __init__(self, atlas_url: str = "http://localhost:21000"):
self.client = Atlas(atlas_url, ("admin", "admin"))
def crawl_datasource(self, ds_config: dict) -> dict:
"""自動爬取資料來源後設資料"""
db_type = ds_config["type"]
if db_type == "mysql":
return self._crawl_mysql(ds_config)
elif db_type == "kafka":
return self._crawl_kafka(ds_config)
elif db_type == "s3":
return self._crawl_s3(ds_config)
elif db_type == "elasticsearch":
return self._crawl_elasticsearch(ds_config)
else:
raise ValueError(f"不支援的資料來源型別: {db_type}")
def infer_lineage_from_sql(self, sql: str) -> list:
"""從SQL語句自動推導血緣關係"""
tree = sqlglot.parse_one(sql, dialect="mysql")
lineage = []
# 提取目標表
target = tree.find(sqlglot.exp.Create).this if tree else None
# 提取源表
sources = []
for tbl in tree.find_all(sqlglot.exp.Table):
if tbl.name != str(target):
sources.append({
"source": tbl.name,
"database": tbl.db if tbl.db else "default"
})
for src in sources:
lineage.append({
"source_table": src["source"],
"source_db": src["database"],
"target_table": str(target.name),
"transform_type": tree.key.lower() if hasattr(tree, 'key') else "unknown"
})
return lineage
def generate_smart_tags(self, table_meta: dict) -> list:
"""基於後設資料自動生成智慧標籤"""
tags = []
columns = table_meta.get("columns", [])
for col in columns:
col_name = col["name"].lower()
if "id" in col_name:
tags.append("標識欄位")
if any(kw in col_name for kw in ["amount", "price", "fee", "money"]):
tags.append("金額欄位-PII")
if any(kw in col_name for kw in ["phone", "mobile", "tel"]):
tags.append("手機號-PII-L3")
if any(kw in col_name for kw in ["email", "mail"]):
tags.append("郵箱-PII-L3")
if any(kw in col_name for kw in ["name", "姓名"]):
tags.append("姓名-PII-L3")
return list(set(tags))
# 演示
collector = ActiveMetadataCollector()
# 模擬採集結果
ds_config = {
"type": "mysql",
"host": "10.0.1.100",
"database": "sales",
"tables": ["orders", "customers"]
}
print("=== 主動後設資料採集演示 ===")
print(f" 資料來源: {ds_config['type']} @ {ds_config['host']}")
print(f" 資料庫: {ds_config['database']}")
# SQL血緣推導
sql = """
CREATE TABLE dwd_sales_order AS
SELECT
o.order_id,
c.customer_name,
o.order_amount,
o.created_at
FROM orders o
JOIN customers c ON o.customer_id = c.customer_id
WHERE o.status = 'paid'
"""
lineage = collector.infer_lineage_from_sql(sql)
print(f"\n 血緣推導: orders + customers → dwd_sales_order")
print(f" 源表: {', '.join(ln['source_table'] for ln in lineage)}")
將後設資料構建為知識圖譜,讓 AI 也能理解資料之間的關係。
-- scripts/knowledge_graph.cypher
-- 用 Neo4j 構建資料資產知識圖譜
// 1. 建立資料來源節點
CREATE (ds1:Datasource {name: '交易庫MySQL', type: 'mysql', owner: '資料平臺組'})
CREATE (ds2:Datasource {name: '日誌叢集Kafka', type: 'kafka', owner: '基礎架構組'})
CREATE (ds3:Datasource {name: '湖倉MinIO', type: 's3', owner: '資料平臺組'})
// 2. 建立資料集節點
CREATE
(ds1)-[:CONTAINS]->(:Dataset {name: 'orders', domain: '交易', freshness: '即時'}),
(ds1)-[:CONTAINS]->(:Dataset {name: 'customers', domain: '客戶', freshness: 'T+1'}),
(ds1)-[:CONTAINS]->(:Dataset {name: 'products', domain: '商品', freshness: 'T+1'}),
(ds2)-[:CONTAINS]->(:Dataset {name: 'user_clicks', domain: '使用者行為', freshness: '即時'}),
(ds3)-[:CONTAINS]->(:Dataset {name: 'dwd_sales_order', domain: '交易', freshness: 'T+0'}),
(ds3)-[:CONTAINS]->(:Dataset {name: 'dws_user_profile', domain: '使用者畫像', freshness: 'T+1'})
// 3. 建立血緣關係(資料從哪裡來、到哪裡去)
MATCH (a:Dataset {name: 'orders'}), (b:Dataset {name: 'dwd_sales_order'})
CREATE (a)-[:FEEDS_INTO]->(b)
MATCH (a:Dataset {name: 'customers'}), (b:Dataset {name: 'dwd_sales_order'})
CREATE (a)-[:FEEDS_INTO]->(b)
MATCH (a:Dataset {name: 'user_clicks'}), (b:Dataset {name: 'dws_user_profile'})
CREATE (a)-[:FEEDS_INTO]->(b)
MATCH (a:Dataset {name: 'dwd_sales_order'}), (b:Dataset {name: 'dws_user_profile'})
CREATE (a)-[:FEEDS_INTO]->(b)
// 4. 知識圖譜查詢演示
// Q: dwd_sales_order 依賴哪些上游表?
// MATCH (upstream)-[:FEEDS_INTO]->(d:Dataset {name:'dwd_sales_order'})
// RETURN upstream.name, upstream.domain
用一個 SQL 查詢多個異構資料來源,無需搬遷資料。
# scripts/federated_query.py
"""聯邦查詢演示:Trino 跨源查詢"""
# Trino 聯邦查詢示例
SQL_FEDERATED_QUERY = """
-- 同時查詢 MySQL 交易表 + Kafka 日誌 + Iceberg 湖表
WITH
-- 源1: MySQL 即時訂單
recent_orders AS (
SELECT customer_id, order_amount, created_at
FROM mysql.sales.orders
WHERE created_at >= DATE('2026-05-01')
),
-- 源2: Iceberg 湖倉中的客戶畫像
customer_profiles AS (
SELECT customer_id, age_group, city, credit_level
FROM iceberg.data_lake.dws_user_profile
),
-- 源3: Kafka 即時行為
user_behavior AS (
SELECT
JSON_EXTRACT_SCALAR(event, '$.customer_id') AS customer_id,
JSON_EXTRACT_SCALAR(event, '$.action') AS action,
count(*) AS action_count
FROM kafka.user_behavior.user_clicks
GROUP BY 1, 2
)
-- 聯邦 JOIN:三源合一,資料不搬遷
SELECT
c.age_group,
c.city,
c.credit_level,
SUM(o.order_amount) AS total_amount,
ub.action,
ub.action_count
FROM recent_orders o
JOIN customer_profiles c ON o.customer_id = c.customer_id
LEFT JOIN (
SELECT customer_id, action, action_count,
ROW_NUMBER() OVER(PARTITION BY customer_id ORDER BY action_count DESC) AS rn
FROM user_behavior
) ub ON o.customer_id = ub.customer_id AND ub.rn = 1
GROUP BY 1, 2, 3, 5, 6
ORDER BY total_amount DESC
LIMIT 100
"""
print("=== 聯邦查詢演示 ===")
print("Trino 同時查詢 MySQL + Kafka + Iceberg(三個異構源)")
print("資料不搬遷,即時 JOIN,秒級返回\n")
print(SQL_FEDERATED_QUERY)
# scripts/data_portal.py
"""資料資產地圖核心——讓業務人員自助發現和使用資料"""
class DataPortal:
"""
資料資產地圖
功能:
- 關鍵詞搜尋資料資產
- 推薦相似資料集
- 一鍵申請訪問許可權
- 檢視資料質量評分
"""
def __init__(self):
self.catalog = self._load_catalog()
def _load_catalog(self):
"""模擬從 Atlas/OpenMetadata 載入目錄"""
return {
"dwd_sales_order": {
"name": "交易訂單明細寬表",
"domain": "交易域",
"freshness": "T+0 日",
"quality_score": 98,
"usage_count_7d": 156,
"columns": [
{"name": "order_id", "description": "訂單ID", "pii": False},
{"name": "customer_name", "description": "客戶姓名", "pii": True, "level": "L3"},
{"name": "order_amount", "description": "訂單金額(元)", "pii": False},
],
"tags": ["P0", "交易", "寬表", "自助取數"],
"owner": "張三(資料平臺組)"
},
"dws_user_profile": {
"name": "使用者畫像彙總表",
"domain": "使用者域",
"freshness": "T+1 日",
"quality_score": 95,
"usage_count_7d": 203,
"tags": ["P0", "使用者畫像", "營銷"],
"owner": "李四(推薦演算法組)"
},
"ads_report_weekly": {
"name": "每週經營分析報表",
"domain": "報表域",
"freshness": "每週一更新",
"quality_score": 100,
"usage_count_7d": 42,
"tags": ["P1", "報表", "管理"],
"owner": "王五(BI組)"
}
}
def search(self, keyword: str) -> list:
"""關鍵詞搜尋資料資產"""
results = []
keyword_lower = keyword.lower()
for id, meta in self.catalog.items():
# 匹配表名、描述、域名、標籤
text = f"{id} {meta['name']} {meta['domain']} {' '.join(meta['tags'])}"
if keyword_lower in text.lower():
results.append({"id": id, **meta})
return results
def get_popular(self, top_n: int = 5) -> list:
"""熱門資料資產"""
sorted_items = sorted(
self.catalog.items(),
key=lambda x: x[1]["usage_count_7d"],
reverse=True
)
return [{"id": id, **meta} for id, meta in sorted_items[:top_n]]
# 演示
portal = DataPortal()
print("=== Data Fabric 資料資產地圖 ===\n")
# 搜尋
print("搜尋「交易」相關資料:")
for r in portal.search("交易"):
print(f" {r['id']:30s} | {r['name']} | 質量: {r['quality_score']}% | 周用量: {r['usage_count_7d']}")
# 熱門榜單
print("\n📊 本週熱門資料 TOP 3:")
for i, r in enumerate(portal.get_popular(3), 1):
print(f" #{i} {r['id']:30s} | {r['name']} | 周訪問 {r['usage_count_7d']}次")
| 元件 | 推薦 | 備選 | 說明 |
|---|---|---|---|
| 後設資料管理 | Apache Atlas | OpenMetadata, DataHub | Atlas 生態最強(Hive/Kafka血緣) |
| 資料目錄 | OpenMetadata | Amundsen, DataHub | UI友好,協作標註 |
| 資料虛擬化 | Trino | Denodo, Starburst | 開源+高效能聯邦查詢 |
| 知識圖譜 | Neo4j | NebulaGraph, JanusGraph | 後設資料關係建模 |
| 自動化編排 | Airflow + dbt | Dagster, Prefect | 按後設資料自動生成 DAG |
| 資料質量 | Great Expectations | Soda, Deequ | 與後設資料聯動 |
階段1(0-1月):後設資料基礎
✓ 部署 Atlas + OpenMetadata
✓ 接入 3 個核心資料來源
✓ 建立資料目錄基礎
階段2(1-2月):血緣與圖譜
✓ 自動血緣採集上線
✓ 構建資料資產知識圖譜
✓ 智慧標籤自動打標
階段3(2-4月):虛擬化與自助
✓ Trino 聯邦查詢上線
✓ 資料資產地圖門戶
✓ 自助取數工作流
階段4(4-6月):智慧編織
✓ AI 推薦相關資料集
✓ 自動治理建議
✓ Fabric → Mesh 過渡評估
Data Fabric(編織) Data Mesh(網格)
───────────────────── ─────────────────
中心化後設資料驅動 去中心化領域自治
統一虛擬化層 各域獨立資料產品
全域性治理 聯邦治理
↓
過渡策略:
1. Fabric 的後設資料層 → Mesh 的聯邦目錄
2. Fabric 的虛擬化層 → Mesh 各域的 Data API
3. Fabric 的治理 → Mesh 的聯邦治理委員會
| 級別 | 特徵 | 評估 |
|---|---|---|
| L1 初始 | 手動 ETL,無後設資料管理 | 原始 |
| L2 可管理 | 有資料目錄,但被動後設資料 | 基礎 |
| L3 主動 | 主動後設資料,自動血緣 | ✅ 當前技能 |
| L4 智慧 | AI 推薦,主動治理建議 | 目標 |
| L5 自適應 | 自動編排,自最佳化 | 願景 |
這個 Skill 資料內容豐富、覆蓋面廣,架構圖清晰、對比分析深入,能幫助理解 Data Fabric 核心理念。但存在明顯短板:文件提到的很多功能指令碼實際不存在,示例程式碼不完整,部分環境配置缺失,導致無法真正執行體驗。適合作為學習參考資料,但實用性打了折扣。