Data Fabric 編織者

👤 莊子十八代技師 📦 v1.0.0 ⭐ 4.3 ⬇️ 225 下載
📊 資料分析 免費

📖 技能介紹


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


Data Fabric 編織者

用後設資料的絲線編織資料生態——讓異構資料來源無縫對話,讓資料消費像開啟 App Store 一樣簡單。

一、Data Fabric 全景圖

                      ┌─────────────────────────────────┐
                      │        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

二、Data Fabric vs 傳統架構

維度 傳統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

四、核心能力詳解

4.1 主動後設資料 — Apache Atlas

主動後設資料不僅記錄"資料在哪裡",更自動感知變化並推送建議。

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)}")

4.2 知識圖譜 — 資料資產地圖

將後設資料構建為知識圖譜,讓 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

4.3 資料虛擬化 — 聯邦查詢

用一個 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)

4.4 資料資產地圖 — 自助服務門戶

# 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 過渡

Data Fabric(編織)           Data Mesh(網格)
─────────────────────         ─────────────────
  中心化後設資料驅動               去中心化領域自治
  統一虛擬化層                   各域獨立資料產品
  全域性治理                       聯邦治理
                                ↓
        過渡策略:
        1. Fabric 的後設資料層 → Mesh 的聯邦目錄
        2. Fabric 的虛擬化層 → Mesh 各域的 Data API
        3. Fabric 的治理 → Mesh 的聯邦治理委員會

八、Gartner Data Fabric 成熟度模型

級別 特徵 評估
L1 初始 手動 ETL,無後設資料管理 原始
L2 可管理 有資料目錄,但被動後設資料 基礎
L3 主動 主動後設資料,自動血緣 ✅ 當前技能
L4 智慧 AI 推薦,主動治理建議 目標
L5 自適應 自動編排,自最佳化 願景

🤖 AI 評測

這個 Skill 資料內容豐富、覆蓋面廣,架構圖清晰、對比分析深入,能幫助理解 Data Fabric 核心理念。但存在明顯短板:文件提到的很多功能指令碼實際不存在,示例程式碼不完整,部分環境配置缺失,導致無法真正執行體驗。適合作為學習參考資料,但實用性打了折扣。

📊 多維度評分

適應性3.9
規範性4.2
有效性4.7
可靠性3.7
可信度4.7

📁 包含檔案 (8 個)

📄 SKILL.md 16.1 KB
📄 references/atlas-openmetadata-comparison.md 3.4 KB
📄 references/data-fabric-architecture.md 9.7 KB
📄 references/data-virtuallization-trino.md 5.8 KB
📄 scripts/docker-compose.yml 3.4 KB
📄 scripts/init_metadata.py 4.3 KB
📄 scripts/load_knowledge_graph.py 2.9 KB
📄 scripts/setup_fabric.sh 1.8 KB