基于OpenClaw構(gòu)建的智能數(shù)據(jù)分析平臺的完整案例
摘要
本文通過一個完整的數(shù)據(jù)分析平臺案例,演示如何使用 OpenClaw 構(gòu)建智能數(shù)據(jù)分析系統(tǒng)。文章涵蓋數(shù)據(jù)采集、數(shù)據(jù)清洗、數(shù)據(jù)分析、可視化展示等核心功能,幫助開發(fā)者掌握 OpenClaw 在數(shù)據(jù)分析場景的應用。通過詳細的系統(tǒng)設計和代碼實現(xiàn),讓讀者了解數(shù)據(jù)分析平臺的完整構(gòu)建過程。
1. 引言 - 數(shù)據(jù)分析平臺概述
1.1 數(shù)據(jù)分析需求
企業(yè)數(shù)據(jù)分析面臨諸多挑戰(zhàn),傳統(tǒng)方案難以滿足現(xiàn)代業(yè)務需求:
| 挑戰(zhàn) | 傳統(tǒng)方案 | OpenClaw方案 |
|---|---|---|
| 數(shù)據(jù)分散 | 手動匯總 | 自動采集整合 |
| 分析門檻高 | 需要專業(yè)分析師 | 自然語言查詢 |
| 響應慢 | 批量處理 | 實時分析 |
| 洞察淺 | 描述性分析 | 預測性分析 |
| 協(xié)作難 | 報告分發(fā) | 智能問答 |
1.2 平臺架構(gòu)設計

1.3 核心功能規(guī)劃
| 功能模塊 | 核心能力 | 技術(shù)實現(xiàn) |
|---|---|---|
| 數(shù)據(jù)采集 | 多源數(shù)據(jù)接入 | 連接器 + 流式采集 |
| 數(shù)據(jù)清洗 | 數(shù)據(jù)質(zhì)量保障 | 規(guī)則引擎 + 異常檢測 |
| 數(shù)據(jù)分析 | 多維度分析 | SQL + ML |
| 智能查詢 | 自然語言交互 | NLP + SQL生成 |
| 可視化 | 圖表展示 | 圖表庫 + 自動推薦 |
2. 數(shù)據(jù)采集模塊
2.1 多源數(shù)據(jù)連接器
from abc import ABC, abstractmethod
from typing import Dict, List, Any, Optional
from dataclasses import dataclass
import pandas as pd
@dataclass
class DataSource:
"""數(shù)據(jù)源配置"""
name: str
type: str
connection: Dict[str, Any]
schema: Optional[Dict] = None
class DataConnector(ABC):
"""數(shù)據(jù)連接器基類"""
@abstractmethod
def connect(self, source: DataSource) -> bool:
"""建立連接"""
pass
@abstractmethod
def fetch(self, query: str) -> pd.DataFrame:
"""獲取數(shù)據(jù)"""
pass
@abstractmethod
def get_schema(self) -> Dict:
"""獲取數(shù)據(jù)結(jié)構(gòu)"""
pass
class DatabaseConnector(DataConnector):
"""數(shù)據(jù)庫連接器"""
def __init__(self):
self.connection = None
self.source = None
def connect(self, source: DataSource) -> bool:
"""建立數(shù)據(jù)庫連接"""
self.source = source
# 根據(jù)數(shù)據(jù)庫類型選擇驅(qū)動
db_type = source.type
if db_type == "mysql":
import pymysql
self.connection = pymysql.connect(
host=source.connection["host"],
port=source.connection.get("port", 3306),
user=source.connection["user"],
password=source.connection["password"],
database=source.connection["database"]
)
elif db_type == "postgresql":
import psycopg2
self.connection = psycopg2.connect(
host=source.connection["host"],
port=source.connection.get("port", 5432),
user=source.connection["user"],
password=source.connection["password"],
database=source.connection["database"]
)
elif db_type == "sqlite":
import sqlite3
self.connection = sqlite3.connect(source.connection["path"])
return self.connection is not None
def fetch(self, query: str) -> pd.DataFrame:
"""執(zhí)行查詢并返回DataFrame"""
if not self.connection:
raise Exception("未建立連接")
return pd.read_sql(query, self.connection)
def get_schema(self) -> Dict:
"""獲取數(shù)據(jù)庫結(jié)構(gòu)"""
if self.source.type == "mysql":
query = f"""
SELECT table_name, column_name, data_type
FROM information_schema.columns
WHERE table_schema = '{self.source.connection["database"]}'
"""
elif self.source.type == "postgresql":
query = f"""
SELECT table_name, column_name, data_type
FROM information_schema.columns
WHERE table_schema = 'public'
"""
else:
return {}
df = self.fetch(query)
# 構(gòu)建schema字典
schema = {}
for _, row in df.iterrows():
table = row["table_name"]
if table not in schema:
schema[table] = {"columns": []}
schema[table]["columns"].append({
"name": row["column_name"],
"type": row["data_type"]
})
return schema
class APIConnector(DataConnector):
"""API連接器"""
def __init__(self):
self.source = None
self.session = None
def connect(self, source: DataSource) -> bool:
"""配置API連接"""
self.source = source
import requests
self.session = requests.Session()
# 設置認證
if "api_key" in source.connection:
self.session.headers["Authorization"] = f"Bearer {source.connection['api_key']}"
return True
def fetch(self, endpoint: str, params: Dict = None) -> pd.DataFrame:
"""獲取API數(shù)據(jù)"""
if not self.session:
raise Exception("未建立連接")
base_url = self.source.connection["base_url"]
url = f"{base_url}/{endpoint}"
response = self.session.get(url, params=params)
if response.status_code == 200:
data = response.json()
# 處理嵌套數(shù)據(jù)
if isinstance(data, list):
return pd.DataFrame(data)
elif isinstance(data, dict):
# 嘗試找到數(shù)據(jù)列表
for key in ["data", "results", "items"]:
if key in data and isinstance(data[key], list):
return pd.DataFrame(data[key])
return pd.DataFrame([data])
raise Exception(f"API請求失敗: {response.status_code}")
def get_schema(self) -> Dict:
"""獲取API數(shù)據(jù)結(jié)構(gòu)"""
# 通過示例數(shù)據(jù)推斷結(jié)構(gòu)
return {}
class FileConnector(DataConnector):
"""文件連接器"""
def __init__(self):
self.source = None
self.data = None
def connect(self, source: DataSource) -> bool:
"""加載文件"""
self.source = source
file_path = source.connection["path"]
file_type = source.connection.get("type", "csv")
if file_type == "csv":
self.data = pd.read_csv(file_path)
elif file_type == "excel":
self.data = pd.read_excel(file_path)
elif file_type == "json":
self.data = pd.read_json(file_path)
elif file_type == "parquet":
self.data = pd.read_parquet(file_path)
else:
return False
return True
def fetch(self, query: str = None) -> pd.DataFrame:
"""獲取數(shù)據(jù)"""
return self.data.copy()
def get_schema(self) -> Dict:
"""獲取數(shù)據(jù)結(jié)構(gòu)"""
if self.data is None:
return {}
return {
"columns": [
{"name": col, "type": str(self.data[col].dtype)}
for col in self.data.columns
]
}
# 使用示例
# 數(shù)據(jù)庫連接
db_source = DataSource(
name="main_db",
type="mysql",
connection={
"host": "localhost",
"user": "root",
"password": "password",
"database": "mydb"
}
)
db_connector = DatabaseConnector()
db_connector.connect(db_source)
df = db_connector.fetch("SELECT * FROM users LIMIT 10")
print(f"獲取 {len(df)} 條記錄")
# API連接
api_source = DataSource(
name="weather_api",
type="api",
connection={
"base_url": "https://api.weather.com",
"api_key": "your_api_key"
}
)
api_connector = APIConnector()
api_connector.connect(api_source)
weather_df = api_connector.fetch("current", {"city": "beijing"})
print(f"天氣數(shù)據(jù): {weather_df.shape}")
# 文件連接
file_source = DataSource(
name="sales_data",
type="file",
connection={
"path": "/data/sales.csv",
"type": "csv"
}
)
file_connector = FileConnector()
file_connector.connect(file_source)
sales_df = file_connector.fetch()
print(f"銷售數(shù)據(jù): {sales_df.shape}")2.2 數(shù)據(jù)采集調(diào)度器
from typing import Dict, List, Callable
from datetime import datetime, timedelta
import threading
import time
@dataclass
class CollectionTask:
"""采集任務"""
id: str
name: str
connector: DataConnector
query: str
schedule: str # cron表達式
callback: Callable
last_run: datetime = None
next_run: datetime = None
status: str = "pending"
class DataCollectionScheduler:
"""數(shù)據(jù)采集調(diào)度器"""
def __init__(self):
self.tasks: Dict[str, CollectionTask] = {}
self.running = False
self.thread = None
def add_task(self, task: CollectionTask):
"""添加采集任務"""
# 計算下次運行時間
task.next_run = self._parse_schedule(task.schedule)
self.tasks[task.id] = task
def remove_task(self, task_id: str):
"""移除采集任務"""
if task_id in self.tasks:
del self.tasks[task_id]
def start(self):
"""啟動調(diào)度器"""
self.running = True
self.thread = threading.Thread(target=self._run_loop, daemon=True)
self.thread.start()
def stop(self):
"""停止調(diào)度器"""
self.running = False
if self.thread:
self.thread.join(timeout=5)
def _run_loop(self):
"""調(diào)度循環(huán)"""
while self.running:
now = datetime.now()
for task in self.tasks.values():
if task.next_run and now >= task.next_run:
self._execute_task(task)
task.last_run = now
task.next_run = self._parse_schedule(task.schedule)
time.sleep(1)
def _execute_task(self, task: CollectionTask):
"""執(zhí)行采集任務"""
try:
task.status = "running"
# 獲取數(shù)據(jù)
df = task.connector.fetch(task.query)
# 調(diào)用回調(diào)
task.callback(df)
task.status = "success"
except Exception as e:
task.status = "failed"
print(f"任務 {task.name} 執(zhí)行失敗: {e}")
def _parse_schedule(self, schedule: str) -> datetime:
"""解析調(diào)度時間"""
# 簡化實現(xiàn):支持簡單格式
# 實際應使用 croniter 庫
if schedule == "every_minute":
return datetime.now() + timedelta(minutes=1)
elif schedule == "every_hour":
return datetime.now() + timedelta(hours=1)
elif schedule == "every_day":
return datetime.now() + timedelta(days=1)
else:
return datetime.now() + timedelta(hours=1)
def get_task_status(self) -> List[Dict]:
"""獲取任務狀態(tài)"""
return [
{
"id": task.id,
"name": task.name,
"status": task.status,
"last_run": task.last_run.isoformat() if task.last_run else None,
"next_run": task.next_run.isoformat() if task.next_run else None
}
for task in self.tasks.values()
]
# 使用示例
scheduler = DataCollectionScheduler()
def process_data(df: pd.DataFrame):
"""處理采集的數(shù)據(jù)"""
print(f"處理 {len(df)} 條數(shù)據(jù)")
# 存儲到數(shù)據(jù)倉庫
# 或觸發(fā)后續(xù)分析
# 添加采集任務
task1 = CollectionTask(
id="task_001",
name="用戶數(shù)據(jù)采集",
connector=db_connector,
query="SELECT * FROM users WHERE created_at > NOW() - INTERVAL 1 HOUR",
schedule="every_hour",
callback=process_data
)
scheduler.add_task(task1)
scheduler.start()3. 數(shù)據(jù)清洗模塊
3.1 數(shù)據(jù)質(zhì)量檢測
from typing import Dict, List, Tuple
import numpy as np
class DataQualityChecker:
"""數(shù)據(jù)質(zhì)量檢測器"""
def __init__(self):
self.rules: List[Dict] = []
def add_rule(self, column: str, rule_type: str, params: Dict = None):
"""添加質(zhì)量規(guī)則"""
self.rules.append({
"column": column,
"type": rule_type,
"params": params or {}
})
def check(self, df: pd.DataFrame) -> Dict:
"""執(zhí)行質(zhì)量檢測"""
results = {
"total_rows": len(df),
"total_columns": len(df.columns),
"issues": [],
"score": 100
}
for rule in self.rules:
column = rule["column"]
rule_type = rule["type"]
params = rule["params"]
if column not in df.columns:
results["issues"].append({
"column": column,
"type": "missing_column",
"message": f"列 {column} 不存在"
})
continue
col_data = df[column]
if rule_type == "not_null":
null_count = col_data.isnull().sum()
if null_count > 0:
results["issues"].append({
"column": column,
"type": "null_values",
"count": null_count,
"percentage": null_count / len(df) * 100
})
elif rule_type == "unique":
dup_count = col_data.duplicated().sum()
if dup_count > 0:
results["issues"].append({
"column": column,
"type": "duplicates",
"count": dup_count
})
elif rule_type == "range":
min_val = params.get("min")
max_val = params.get("max")
if min_val is not None:
below_min = (col_data < min_val).sum()
if below_min > 0:
results["issues"].append({
"column": column,
"type": "below_min",
"count": below_min,
"min": min_val
})
if max_val is not None:
above_max = (col_data > max_val).sum()
if above_max > 0:
results["issues"].append({
"column": column,
"type": "above_max",
"count": above_max,
"max": max_val
})
elif rule_type == "pattern":
pattern = params.get("pattern")
if pattern:
import re
invalid = ~col_data.astype(str).str.match(pattern, na=False)
invalid_count = invalid.sum()
if invalid_count > 0:
results["issues"].append({
"column": column,
"type": "pattern_mismatch",
"count": invalid_count,
"pattern": pattern
})
# 計算質(zhì)量分數(shù)
if results["issues"]:
issue_penalty = sum(
issue.get("percentage", 5)
for issue in results["issues"]
)
results["score"] = max(0, 100 - issue_penalty)
return results
def get_profile(self, df: pd.DataFrame) -> Dict:
"""獲取數(shù)據(jù)概要"""
profile = {
"row_count": len(df),
"column_count": len(df.columns),
"memory_usage": df.memory_usage(deep=True).sum(),
"columns": {}
}
for col in df.columns:
col_data = df[col]
col_profile = {
"dtype": str(col_data.dtype),
"null_count": col_data.isnull().sum(),
"null_percentage": col_data.isnull().sum() / len(df) * 100,
"unique_count": col_data.nunique()
}
# 數(shù)值類型統(tǒng)計
if col_data.dtype in ["int64", "float64"]:
col_profile.update({
"min": col_data.min(),
"max": col_data.max(),
"mean": col_data.mean(),
"median": col_data.median(),
"std": col_data.std()
})
# 字符串類型統(tǒng)計
elif col_data.dtype == "object":
col_profile.update({
"min_length": col_data.astype(str).str.len().min(),
"max_length": col_data.astype(str).str.len().max(),
"avg_length": col_data.astype(str).str.len().mean()
})
profile["columns"][col] = col_profile
return profile
# 使用示例
checker = DataQualityChecker()
# 添加質(zhì)量規(guī)則
checker.add_rule("user_id", "not_null")
checker.add_rule("user_id", "unique")
checker.add_rule("age", "range", {"min": 0, "max": 150})
checker.add_rule("email", "pattern", {"pattern": r"^[\w\.-]+@[\w\.-]+\.\w+$"})
# 執(zhí)行檢測
df = pd.DataFrame({
"user_id": [1, 2, 3, None, 5],
"age": [25, 30, -5, 40, 200],
"email": ["a@b.com", "invalid", "c@d.com", "e@f.com", "g@h.com"]
})
results = checker.check(df)
print(f"質(zhì)量分數(shù): {results['score']}")
print(f"問題數(shù): {len(results['issues'])}")
# 獲取數(shù)據(jù)概要
profile = checker.get_profile(df)
print(f"數(shù)據(jù)概要: {profile['row_count']} 行, {profile['column_count']} 列")3.2 數(shù)據(jù)清洗處理器
from typing import Dict, List, Callable
class DataCleaner:
"""數(shù)據(jù)清洗處理器"""
def __init__(self):
self.steps: List[Dict] = []
def add_step(self, name: str, processor: Callable, params: Dict = None):
"""添加清洗步驟"""
self.steps.append({
"name": name,
"processor": processor,
"params": params or {}
})
def clean(self, df: pd.DataFrame) -> Tuple[pd.DataFrame, Dict]:
"""執(zhí)行清洗"""
original_count = len(df)
report = {
"original_rows": original_count,
"steps": []
}
result_df = df.copy()
for step in self.steps:
before_count = len(result_df)
result_df = step["processor"](result_df, **step["params"])
after_count = len(result_df)
report["steps"].append({
"name": step["name"],
"rows_before": before_count,
"rows_after": after_count,
"rows_removed": before_count - after_count
})
report["final_rows"] = len(result_df)
report["rows_removed"] = original_count - len(result_df)
return result_df, report
# 預定義清洗處理器
def remove_duplicates(df: pd.DataFrame, subset: List[str] = None) -> pd.DataFrame:
"""去除重復"""
return df.drop_duplicates(subset=subset)
def fill_missing(df: pd.DataFrame, columns: Dict[str, Any] = None) -> pd.DataFrame:
"""填充缺失值"""
result = df.copy()
for col, value in (columns or {}).items():
if col in result.columns:
result[col] = result[col].fillna(value)
return result
def remove_outliers(df: pd.DataFrame, column: str, method: str = "iqr", threshold: float = 1.5) -> pd.DataFrame:
"""去除異常值"""
if column not in df.columns:
return df
if method == "iqr":
Q1 = df[column].quantile(0.25)
Q3 = df[column].quantile(0.75)
IQR = Q3 - Q1
lower = Q1 - threshold * IQR
upper = Q3 + threshold * IQR
return df[(df[column] >= lower) & (df[column] <= upper)]
elif method == "zscore":
from scipy import stats
z_scores = stats.zscore(df[column].dropna())
return df[abs(z_scores) <= threshold]
return df
def standardize_text(df: pd.DataFrame, column: str, lowercase: bool = True, strip: bool = True) -> pd.DataFrame:
"""標準化文本"""
result = df.copy()
if column in result.columns:
if lowercase:
result[column] = result[column].astype(str).str.lower()
if strip:
result[column] = result[column].astype(str).str.strip()
return result
def convert_types(df: pd.DataFrame, columns: Dict[str, str]) -> pd.DataFrame:
"""轉(zhuǎn)換數(shù)據(jù)類型"""
result = df.copy()
for col, dtype in columns.items():
if col in result.columns:
try:
result[col] = result[col].astype(dtype)
except Exception as e:
print(f"轉(zhuǎn)換 {col} 失敗: {e}")
return result
# 使用示例
cleaner = DataCleaner()
# 添加清洗步驟
cleaner.add_step("去除重復", remove_duplicates, {"subset": ["user_id"]})
cleaner.add_step("填充缺失", fill_missing, {"columns": {"age": 0, "name": "未知"}})
cleaner.add_step("去除異常值", remove_outliers, {"column": "age", "method": "iqr"})
cleaner.add_step("標準化文本", standardize_text, {"column": "email", "lowercase": True})
cleaner.add_step("類型轉(zhuǎn)換", convert_types, {"columns": {"age": "int64", "created_at": "datetime64"}})
# 執(zhí)行清洗
cleaned_df, report = cleaner.clean(df)
print(f"清洗報告: {report}")4. 數(shù)據(jù)分析引擎
4.1 統(tǒng)計分析器
from typing import Dict, List, Optional
from scipy import stats
import numpy as np
class StatisticalAnalyzer:
"""統(tǒng)計分析器"""
def __init__(self, df: pd.DataFrame):
self.df = df
def descriptive_stats(self, columns: List[str] = None) -> Dict:
"""描述性統(tǒng)計"""
if columns is None:
columns = self.df.select_dtypes(include=[np.number]).columns.tolist()
result = {}
for col in columns:
if col not in self.df.columns:
continue
data = self.df[col].dropna()
result[col] = {
"count": len(data),
"mean": data.mean(),
"std": data.std(),
"min": data.min(),
"q1": data.quantile(0.25),
"median": data.median(),
"q3": data.quantile(0.75),
"max": data.max(),
"skewness": data.skew(),
"kurtosis": data.kurtosis()
}
return result
def correlation_analysis(self, method: str = "pearson") -> pd.DataFrame:
"""相關(guān)性分析"""
numeric_df = self.df.select_dtypes(include=[np.number])
return numeric_df.corr(method=method)
def hypothesis_test(self, column1: str, column2: str, test_type: str = "ttest") -> Dict:
"""假設檢驗"""
data1 = self.df[column1].dropna()
data2 = self.df[column2].dropna()
if test_type == "ttest":
statistic, pvalue = stats.ttest_ind(data1, data2)
return {
"test": "t-test",
"statistic": statistic,
"p_value": pvalue,
"significant": pvalue < 0.05
}
elif test_type == "mannwhitney":
statistic, pvalue = stats.mannwhitneyu(data1, data2)
return {
"test": "Mann-Whitney U",
"statistic": statistic,
"p_value": pvalue,
"significant": pvalue < 0.05
}
elif test_type == "chi2":
contingency = pd.crosstab(self.df[column1], self.df[column2])
statistic, pvalue, dof, expected = stats.chi2_contingency(contingency)
return {
"test": "Chi-square",
"statistic": statistic,
"p_value": pvalue,
"dof": dof,
"significant": pvalue < 0.05
}
return {}
def anova(self, group_column: str, value_column: str) -> Dict:
"""方差分析"""
groups = self.df.groupby(group_column)[value_column]
group_data = [group.dropna().values for name, group in groups]
statistic, pvalue = stats.f_oneway(*group_data)
return {
"test": "ANOVA",
"statistic": statistic,
"p_value": pvalue,
"significant": pvalue < 0.05,
"groups": len(group_data)
}
def time_series_analysis(self, date_column: str, value_column: str, freq: str = "D") -> Dict:
"""時間序列分析"""
df = self.df.copy()
df[date_column] = pd.to_datetime(df[date_column])
df = df.set_index(date_column)
# 重采樣
resampled = df[value_column].resample(freq)
result = {
"daily_stats": {
"mean": resampled.mean().to_dict(),
"sum": resampled.sum().to_dict(),
"count": resampled.count().to_dict()
}
}
# 趨勢分析
from scipy.signal import detrend
values = df[value_column].values
trend = detrend(values)
result["detrended"] = trend.tolist()
return result
# 使用示例
analyzer = StatisticalAnalyzer(sales_df)
# 描述性統(tǒng)計
stats_result = analyzer.descriptive_stats(["price", "quantity", "total"])
print(f"統(tǒng)計結(jié)果: {stats_result}")
# 相關(guān)性分析
corr = analyzer.correlation_analysis()
print(f"相關(guān)性矩陣:\n{corr}")
# 假設檢驗
test_result = analyzer.hypothesis_test("group_a", "group_b", "ttest")
print(f"檢驗結(jié)果: {test_result}")4.2 自然語言查詢
from typing import Dict, List, Optional
import re
class NaturalLanguageQuery:
"""自然語言查詢處理器"""
def __init__(self, schema: Dict):
self.schema = schema
self.query_templates = self._build_templates()
def _build_templates(self) -> List[Dict]:
"""構(gòu)建查詢模板"""
return [
{
"pattern": r"(.+)的平均值",
"sql_template": "SELECT AVG({column}) FROM {table}",
"type": "aggregation"
},
{
"pattern": r"(.+)的總和",
"sql_template": "SELECT SUM({column}) FROM {table}",
"type": "aggregation"
},
{
"pattern": r"(.+)的最大值",
"sql_template": "SELECT MAX({column}) FROM {table}",
"type": "aggregation"
},
{
"pattern": r"(.+)的最小值",
"sql_template": "SELECT MIN({column}) FROM {table}",
"type": "aggregation"
},
{
"pattern": r"按(.+)分組統(tǒng)計(.+)",
"sql_template": "SELECT {group_column}, COUNT(*) FROM {table} GROUP BY {group_column}",
"type": "grouping"
},
{
"pattern": r"(.+)前(\d+)名",
"sql_template": "SELECT * FROM {table} ORDER BY {column} DESC LIMIT {limit}",
"type": "ranking"
}
]
def parse(self, question: str) -> Dict:
"""解析自然語言問題"""
result = {
"question": question,
"sql": None,
"type": None,
"confidence": 0
}
for template in self.query_templates:
match = re.search(template["pattern"], question)
if match:
result["type"] = template["type"]
# 提取參數(shù)
if template["type"] == "aggregation":
column_name = match.group(1)
column = self._find_column(column_name)
table = self._find_table(column)
if column and table:
result["sql"] = template["sql_template"].format(
column=column,
table=table
)
result["confidence"] = 0.8
elif template["type"] == "grouping":
group_col_name = match.group(1)
value_col_name = match.group(2)
group_column = self._find_column(group_col_name)
table = self._find_table(group_column)
if group_column and table:
result["sql"] = template["sql_template"].format(
group_column=group_column,
table=table
)
result["confidence"] = 0.7
elif template["type"] == "ranking":
column_name = match.group(1)
limit = match.group(2)
column = self._find_column(column_name)
table = self._find_table(column)
if column and table:
result["sql"] = template["sql_template"].format(
column=column,
table=table,
limit=limit
)
result["confidence"] = 0.8
break
return result
def _find_column(self, name: str) -> Optional[str]:
"""查找匹配的列名"""
name_lower = name.lower()
for table_name, table_info in self.schema.items():
for column in table_info.get("columns", []):
if name_lower in column["name"].lower():
return column["name"]
return None
def _find_table(self, column: str) -> Optional[str]:
"""查找列所在的表"""
for table_name, table_info in self.schema.items():
for col in table_info.get("columns", []):
if col["name"] == column:
return table_name
return None
def execute(self, question: str, connector: DataConnector) -> pd.DataFrame:
"""執(zhí)行自然語言查詢"""
parsed = self.parse(question)
if parsed["sql"]:
return connector.fetch(parsed["sql"])
return pd.DataFrame()
# 使用示例
schema = {
"sales": {
"columns": [
{"name": "product", "type": "string"},
{"name": "price", "type": "float"},
{"name": "quantity", "type": "int"},
{"name": "region", "type": "string"}
]
}
}
nlq = NaturalLanguageQuery(schema)
# 解析問題
result = nlq.parse("價格的平均值")
print(f"SQL: {result['sql']}")
print(f"置信度: {result['confidence']}")
# 執(zhí)行查詢
# df = nlq.execute("銷售額前10名", db_connector)5. 可視化引擎
5.1 圖表生成器
from typing import Dict, List, Optional
import matplotlib.pyplot as plt
import seaborn as sns
class ChartGenerator:
"""圖表生成器"""
def __init__(self, style: str = "seaborn"):
plt.style.use(style)
self.figures: List[plt.Figure] = []
def bar_chart(self, df: pd.DataFrame, x: str, y: str, title: str = None) -> plt.Figure:
"""柱狀圖"""
fig, ax = plt.subplots(figsize=(10, 6))
df.plot.bar(x=x, y=y, ax=ax)
if title:
ax.set_title(title)
ax.set_xlabel(x)
ax.set_ylabel(y)
plt.tight_layout()
self.figures.append(fig)
return fig
def line_chart(self, df: pd.DataFrame, x: str, y: str, title: str = None) -> plt.Figure:
"""折線圖"""
fig, ax = plt.subplots(figsize=(12, 6))
df.plot.line(x=x, y=y, ax=ax)
if title:
ax.set_title(title)
plt.tight_layout()
self.figures.append(fig)
return fig
def pie_chart(self, df: pd.DataFrame, values: str, labels: str, title: str = None) -> plt.Figure:
"""餅圖"""
fig, ax = plt.subplots(figsize=(8, 8))
ax.pie(df[values], labels=df[labels], autopct='%1.1f%%')
if title:
ax.set_title(title)
self.figures.append(fig)
return fig
def scatter_plot(self, df: pd.DataFrame, x: str, y: str, hue: str = None, title: str = None) -> plt.Figure:
"""散點圖"""
fig, ax = plt.subplots(figsize=(10, 8))
if hue:
for category in df[hue].unique():
subset = df[df[hue] == category]
ax.scatter(subset[x], subset[y], label=category, alpha=0.6)
ax.legend()
else:
ax.scatter(df[x], df[y], alpha=0.6)
ax.set_xlabel(x)
ax.set_ylabel(y)
if title:
ax.set_title(title)
plt.tight_layout()
self.figures.append(fig)
return fig
def heatmap(self, df: pd.DataFrame, title: str = None) -> plt.Figure:
"""熱力圖"""
fig, ax = plt.subplots(figsize=(10, 8))
sns.heatmap(df, annot=True, fmt=".2f", cmap="coolwarm", ax=ax)
if title:
ax.set_title(title)
plt.tight_layout()
self.figures.append(fig)
return fig
def histogram(self, df: pd.DataFrame, column: str, bins: int = 30, title: str = None) -> plt.Figure:
"""直方圖"""
fig, ax = plt.subplots(figsize=(10, 6))
ax.hist(df[column], bins=bins, edgecolor='black')
ax.set_xlabel(column)
ax.set_ylabel('Frequency')
if title:
ax.set_title(title)
plt.tight_layout()
self.figures.append(fig)
return fig
def box_plot(self, df: pd.DataFrame, x: str, y: str, title: str = None) -> plt.Figure:
"""箱線圖"""
fig, ax = plt.subplots(figsize=(10, 6))
df.boxplot(column=y, by=x, ax=ax)
if title:
ax.set_title(title)
plt.tight_layout()
self.figures.append(fig)
return fig
def save_all(self, directory: str, prefix: str = "chart"):
"""保存所有圖表"""
import os
os.makedirs(directory, exist_ok=True)
for i, fig in enumerate(self.figures):
path = os.path.join(directory, f"{prefix}_{i+1}.png")
fig.savefig(path, dpi=150)
print(f"保存圖表: {path}")
def close_all(self):
"""關(guān)閉所有圖表"""
for fig in self.figures:
plt.close(fig)
self.figures.clear()
# 使用示例
chart_gen = ChartGenerator()
# 生成圖表
chart_gen.bar_chart(sales_df, x="product", y="sales", title="產(chǎn)品銷售額")
chart_gen.line_chart(time_df, x="date", y="revenue", title="收入趨勢")
chart_gen.pie_chart(category_df, values="amount", labels="category", title="類別占比")
chart_gen.scatter_plot(sales_df, x="price", y="quantity", hue="region", title="價格與銷量關(guān)系")
# 保存圖表
chart_gen.save_all("/output/charts", "analysis")6. 最佳實踐
6.1 平臺設計原則
| 原則 | 說明 | 實踐 |
|---|---|---|
| 易用性 | 降低使用門檻 | 自然語言查詢 |
| 可擴展 | 支持新數(shù)據(jù)源 | 插件式連接器 |
| 高性能 | 快速響應 | 緩存 + 索引 |
| 安全性 | 數(shù)據(jù)保護 | 權(quán)限控制 |
6.2 常見問題
| 問題 | 原因 | 解決方案 |
|---|---|---|
| 查詢慢 | 數(shù)據(jù)量大 | 分區(qū) + 索引 |
| 結(jié)果不準 | 數(shù)據(jù)質(zhì)量差 | 數(shù)據(jù)清洗 |
| 圖表亂碼 | 編碼問題 | 設置字體 |
7. 總結(jié)
7.1 核心要點
本文通過完整的數(shù)據(jù)分析平臺案例,展示了 OpenClaw 在數(shù)據(jù)分析場景的應用:
| 模塊 | 核心功能 | 技術(shù)要點 |
|---|---|---|
| 數(shù)據(jù)采集 | 多源接入 | 連接器 + 調(diào)度 |
| 數(shù)據(jù)清洗 | 質(zhì)量保障 | 規(guī)則 + 處理器 |
| 數(shù)據(jù)分析 | 多維分析 | 統(tǒng)計 + ML |
| 智能查詢 | 自然語言 | NLP + SQL |
| 可視化 | 圖表展示 | 自動推薦 |
以上就是基于OpenClaw構(gòu)建的智能數(shù)據(jù)分析平臺的完整案例的詳細內(nèi)容,更多關(guān)于OpenClaw構(gòu)建的智能數(shù)據(jù)分析平臺的資料請關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章

Windows/Mac/Linux下OpenClaw代理配置差異與實操指南
如果你在多臺設備上部署 OpenClaw,你一定遇到過這種跨平臺水土不服的問題,不同操作系統(tǒng)的終端環(huán)境不同,命令語法不同,甚至環(huán)境變量的生效方式也不同,今天這篇文章,就幫2026-07-01
Docker 部署是 OpenClaw 從開發(fā)環(huán)境走向生產(chǎn)環(huán)境的關(guān)鍵一步,本文從 Docker 基礎概念出發(fā),深入講解 OpenClaw 的 Docker 鏡像構(gòu)建、多階段優(yōu)化、環(huán)境變量注入、數(shù)據(jù)持久化、2026-06-30
OpenClaw中3個提效設置實戰(zhàn):自動快模式、自適應思考、定時工作流
OpenClaw 是自托管的 AI Gateway,核心能力是把聊天應用(Discord、Telegram、微信、飛書等)跟各種 AI 模型連接起來,本文為大家分享了3個OpenClaw提效設置技巧,感興趣的2026-06-30
使用OpenClaw Browser實現(xiàn)表單自動化的方法
本文通過實際案例演示 OpenClaw Browser 的表單自動化能力,從簡單登錄表單到復雜多步驟表單,全面解析表單自動化的實現(xiàn)技巧,涵蓋表單分析、智能填寫、驗證處理、錯誤恢復等2026-06-29
OpenClaw從單機Docker部署遷移到Kubernetes集群的完整方案
當你的 OpenClaw 從單機走向集群,Kubernetes 是繞不開的選擇,本文從 K8s 核心概念出發(fā),系統(tǒng)講解 OpenClaw 的 K8s 部署架構(gòu),需要的朋友可以參考下2026-06-29
本文詳細介紹 OpenClaw Canvas 的截圖功能,從基本截圖、全頁面捕獲、元素截圖到圖像處理,全面解析如何通過 Canvas 實現(xiàn)靈活的頁面捕獲,通過實際案例演示報告生成、內(nèi)容存2026-06-28
OpenClaw中間件請求攔截、轉(zhuǎn)換與增強的完整指南
中間件是 OpenClaw 處理鏈路中最靈活的一環(huán),本文從中間件的設計哲學出發(fā),系統(tǒng)講解中間件的三種模式(前置、后置、環(huán)繞)、洋蔥模型執(zhí)行鏈、請求/響應變換機制,以及流式消2026-06-28
Ubuntu從零部署OpenClaw全過程(本地模型+DeepSeek)
OpenClaw 給是一個開源、可自托管的 AI 助手平臺,原生支持 Ollama 本地模型和 DeepSeek 等云端 API,讓你在隱私與性能之間自由切換,本文記錄了我在 Ubuntu 上從零部署 Ope2026-06-26
OpenClaw Token節(jié)省指南:Token消耗如何直降 90%?
QMD(Quantum Memory Database) 的本地語義檢索引擎正在改變這個局面,它用“先檢索、后推理”的思路,把 Token 消耗砍掉了 90% 以上,這篇文章就來深入拆解 QMD 的技術(shù)原2026-06-23
很多人第一次聽到 OpenClaw Skill,會把它理解成“插件”, 這個理解只對了一半, 插件通常給 Agent 增加新的能力,比如新的工具、新的消息渠道、新的模型 Provider,下面我2026-06-23











