文章
构建DuckDB 数据治理平台:算子引擎、安全机制与流水线设计
一、写在前面
在数据治理的日常工作中,我们经常需要处理各种非结构化数据——从 PDF 中提取的文本、爬虫抓取的网页内容、API 返回的 JSON 数据。这些数据往往包含空值、重复、格式混乱等问题,需要一个高效、灵活的数据清洗方案。
传统做法是写 Python 脚本,用 Pandas 进行清洗。但随着数据量增长(几十 GB 甚至更大),Pandas 的内存瓶颈开始显现——数据必须全部加载到内存中,稍有不慎就会 OOM。
于是我将目光投向了 DuckDB:一个嵌入式、列式存储、支持 SQL 的 OLAP 引擎。它可以直接读取 CSV、Parquet 文件,支持流式处理,数据量远超内存也能跑。
基于 DuckDB,我构建了一个完整的 数据治理平台,核心功能包括:
- 算子引擎:用户可注册自定义 Python UDF,在 SQL 中直接调用
- 治理流水线:按步骤执行 SQL,视图透传,支持链式调用
- 安全机制:限制用户代码只能导入白名单模块,拦截危险函数
- 配置驱动:MySQL 存储项目配置,支持前端在线编排
本文将分享这个项目的设计思路与核心实现。
二、为什么选择 DuckDB?
| 对比项 | Pandas | DuckDB |
|---|---|---|
| 内存使用 | 全量加载到内存 | 流式处理,可处理 > 内存的数据 |
| 数据源 | CSV、Excel | CSV、JSON、Parquet、Excel、数据库 |
| 操作方式 | Python API | SQL(支持标准 SQL + 窗口函数) |
| 性能 | 万行以内良好 | 百万行秒级响应 |
对数据治理工程师来说,DuckDB 最大的价值是:可以直接用 SQL 操作本地文件,无需启动任何数据库服务。
SELECT * FROM read_csv('data/*.csv') WHERE content IS NOT NULL三、整体架构
┌─────────────────────────────────────────────────────────────┐
│ 前端(可视化管理界面) │
│ 编排治理步骤 / 注册算子 / 启动任务 │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ MySQL 配置库 │
│ project / operator / step / project_operator │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 治理引擎(Python) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │算子注册器│ │流水线引擎│ │安全模块 │ │
│ └──────────┘ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 数据源(文件系统) │
│ CSV / JSON / Parquet / Excel │
└─────────────────────────────────────────────────────────────┘四、核心能力详解
4.1 算子引擎:让 SQL 拥有 Python 的能力
DuckDB 支持 UDF(用户定义函数),但需要提前注册 Python 函数。我们将其封装为 算子引擎:
核心设计:
- 算子注册中心:单例模式,全局唯一,同名禁止覆盖
- 动态加载:从 MySQL 读取算子代码,通过
exec动态编译并注册到 DuckDB - 版本管理:基于代码内容自动生成版本号
# 使用示例
engine = DataGovernanceEngine()
# 注册一个清洗文本的算子
engine.register_operator(
op_name="clean_text",
code="""
def clean_text(text: str) -> str:
if text is None:
return ''
return ' '.join(text.split())
""",
creator="admin"
)
# 在 SQL 中直接调用
result = engine.sql("SELECT clean_text(content) FROM raw_data")4.2 治理流水线:分步执行,视图透传
数据治理往往需要多步操作:加载 → 清洗 → 过滤 → 去重 → 聚合。我们将每一步抽象为流水线步骤,前一步的输出作为后一步的输入。
核心设计:
- 使用
CREATE OR REPLACE VIEW保存中间结果 - SQL 模板支持
{view}变量(引用上一步输出) - 支持通配符批量加载文件
# 流水线示例
pipeline = engine.create_pipeline()
pipeline.set_file_path("data/*.csv")
pipeline.add_step("加载", "CREATE OR REPLACE VIEW raw AS SELECT * FROM read_csv('{{file_path}}')", output_view="raw")
pipeline.add_step("清洗", "CREATE OR REPLACE VIEW cleaned AS SELECT *, clean_text(content) AS clean FROM {view}", output_view="cleaned")
pipeline.add_step("去重", "CREATE OR REPLACE VIEW deduped AS SELECT * FROM {view} QUALIFY ROW_NUMBER() OVER (PARTITION BY clean) = 1", output_view="deduped")
pipeline.add_step("报告", "SELECT source, COUNT(*) FROM {view} GROUP BY source")
result = pipeline.run()4.3 安全机制:让用户代码受控运行
用户在线编写的 Python 代码通过 exec 执行,存在安全风险。我们通过两层防护确保安全:
第一层:AST 静态分析
在 exec 之前,解析代码的 AST(抽象语法树),检测是否包含危险函数引用(eval、open、exec、__import__ 等)。如果检测到,直接拒绝注册,无需实际执行代码。
def _has_dangerous_reference(code: str) -> bool:
dangerous_names = {'eval', 'exec', 'compile', 'open', 'input', '__import__'}
tree = ast.parse(code)
for node in ast.walk(tree):
if isinstance(node, ast.Name) and node.id in dangerous_names:
return True
if isinstance(node, ast.Call) and isinstance(node.func, ast.Name):
if node.func.id in dangerous_names:
return True
return False第二层:白名单导入控制
劫持 __import__,只允许导入白名单中的模块。白名单分为两类:
- 标准库(安全子集):
math、re、json、datetime、pandas等 - 额外允许的第三方库:用户通过
extra_libraries参数显式声明
class SecureImporter:
def _safe_import(self, name, *args, **kwargs):
top_level = name.split('.')[0]
if top_level not in self._allowed_modules:
raise SecurityError(f"禁止导入模块: {name}")
return self._original_import(name, *args, **kwargs)效果对比:
# ❌ 被拦截:导入危险模块
import os # SecurityError: 禁止导入模块: 'os'
# ❌ 被拦截:调用危险函数
def evil():
eval("__import__('os').system('rm -rf /')") # AST 检测到 eval 引用
# ✅ 被允许:传递 extra_libraries
engine.register_operator("my_func", code, extra_libraries=['pandas', 'numpy'])4.4 配置驱动:MySQL 存储,前端可编排
将治理流程抽象为 项目(Project)→ 步骤(Step)→ 算子(Operator) 三层结构,存储在 MySQL 中。前端可以可视化管理:
核心表结构:
governance_project:项目信息governance_operator:算子定义(Python 代码)governance_step:治理步骤(SQL 模板 + 执行顺序)
一键执行:
result = engine.run_project_from_mysql(
db_url="mysql://user:pass@host:3306/governance_platform",
project_name="数据质量治理",
file_path="data/*.csv"
)五、实际效果:100 条 CSV 数据治理
原始数据(100 行):
- 8% NULL 内容(解析失败)
- 7% 空字符串或制表符
- 7% 重复内容
治理流程:加载 → 清洗 → 过滤(长度 >= 10) → 去重(保留质量分最高) → 聚合报告
治理前后对比:
| 指标 | 治理前 | 治理后 |
|---|---|---|
| 总行数 | 100 | 81 |
| NULL 内容 | 15 行 | 0 |
| 空内容 | 4 行 | 0 |
| 去重后内容 | 78 行 | 81(保留唯一) |
治理报告示例:
source category total avg_len high_quality
api finance 6 28.67 2
api tech 3 32.33 1
excel education 7 30.43 1
excel medical 5 30.60 3六、技术选型与依赖
| 组件 | 用途 | 版本 |
|---|---|---|
| DuckDB | 嵌入式 SQL 引擎,数据治理核心 | ≥ 0.9.0 |
| SQLAlchemy | MySQL 连接池管理 | ≥ 2.0.0 |
| PyMySQL | MySQL 驱动 | ≥ 1.0.0 |
| Pandas | 结果展示 | ≥ 2.0.0 |
| Python | 开发语言 | ≥ 3.9 |
项目代码示例

core/exceptions.py
"""自定义异常模块"""
class OperatorExistsError(Exception):
"""算子已存在异常(禁止覆盖)"""
def __init__(self, op_name: str, creator: str, created_at: str):
self.op_name = op_name
self.creator = creator
self.created_at = created_at
super().__init__(
f"算子 '{op_name}' 已存在,由 {creator} 于 {created_at} 创建,禁止覆盖"
)
class OperatorNotFoundError(Exception):
"""算子不存在异常"""
def __init__(self, op_name: str):
self.op_name = op_name
super().__init__(f"算子 '{op_name}' 不存在")
class CodeCompileError(Exception):
"""代码编译失败异常"""
pass
class StepExecutionError(Exception):
"""步骤执行失败异常"""
def __init__(self, step_name: str, sql: str, error: str):
self.step_name = step_name
self.sql = sql
self.error = error
super().__init__(f"步骤 '{step_name}' 执行失败: {error}")
core/security.py
"""
算子安全模块:限制用户代码只能导入白名单中的模块
职责:
1. 劫持 __import__ 拦截所有导入请求
2. 白名单控制:标准库(安全子集)+ 用户允许的第三方库
3. 危险内置函数(eval, open 等)由 AST 静态分析拦截(在 dynamic_loader 中)
"""
import sys
import builtins
from typing import Set, List, Optional
# ================================================================
# 1. 标准库白名单(安全子集,不含危险模块)
# ================================================================
def _get_stdlib_modules() -> Set[str]:
"""获取安全的标准库模块名称列表"""
if hasattr(sys, 'stdlib_module_names'):
stdlib = set(sys.stdlib_module_names)
else:
stdlib = set(sys.builtin_module_names)
common_stdlib = {
'math', 're', 'json', 'datetime', 'time', 'random',
'string', 'collections', 'itertools', 'functools',
'typing', 'hashlib', 'base64', 'binascii', 'csv', 'io',
'pathlib', 'glob', 'shutil', 'copy', 'pprint', 'textwrap',
'unicodedata', 'calendar', 'decimal', 'fractions',
'statistics', 'numbers', 'operator', 'inspect', 'pickle',
'warnings', 'contextlib', 'tempfile', 'enum',
}
stdlib.update(common_stdlib)
# 移除危险模块
dangerous_stdlib = {
'os', 'socket', 'subprocess', 'multiprocessing', 'threading',
'ctypes', 'cffi', 'distutils', 'venv', 'pdb',
'code', 'codeop', 'compileall', 'py_compile',
'importlib', 'imp', 'sysconfig', 'site',
}
return stdlib - dangerous_stdlib
_STANDARD_LIBRARIES: Set[str] = _get_stdlib_modules()
# ================================================================
# 2. 安全导入器
# ================================================================
class SecurityError(Exception):
"""安全违规异常"""
def __init__(self, module_name: str, allowed: Set[str]):
self.module_name = module_name
self.allowed = allowed
super().__init__(
f"禁止导入模块: '{module_name}'。"
f"只允许: {', '.join(sorted(allowed)[:10])}..."
)
class SecureImporter:
def __init__(self, extra_libraries: Optional[List[str]] = None):
self._allowed_modules = set(_STANDARD_LIBRARIES)
if extra_libraries:
self._allowed_modules.update(extra_libraries)
self._original_import = builtins.__import__
def _safe_import(self, name: str, globals_dict: Optional[dict] = None,
locals_dict: Optional[dict] = None, fromlist: Optional[list] = None,
level: int = 0):
top_level = name.split('.')[0]
if top_level not in self._allowed_modules:
# 检查是否是已允许模块的子模块
if not any(name.startswith(f"{allowed}.") for allowed in self._allowed_modules):
raise SecurityError(name, self._allowed_modules)
return self._original_import(name, globals_dict, locals_dict, fromlist, level)
def create_safe_builtins(self) -> dict:
"""创建安全的内置函数字典(不含危险函数)"""
safe_builtins = {
'len': len, 'str': str, 'int': int, 'float': float,
'bool': bool, 'list': list, 'dict': dict, 'tuple': tuple,
'set': set, 'frozenset': frozenset, 'type': type,
'range': range, 'enumerate': enumerate, 'zip': zip,
'map': map, 'filter': filter, 'sorted': sorted,
'reversed': reversed, 'sum': sum, 'min': min, 'max': max,
'any': any, 'all': all, 'abs': abs, 'round': round,
'pow': pow, 'divmod': divmod, 'hex': hex, 'oct': oct, 'bin': bin,
'ord': ord, 'chr': chr, 'repr': repr, 'format': format,
'hash': hash, 'id': id, 'iter': iter, 'next': next,
'isinstance': isinstance, 'issubclass': issubclass,
'hasattr': hasattr, 'getattr': getattr, 'setattr': setattr,
'delattr': delattr, 'callable': callable, 'dir': dir,
'help': help, 'print': print,
# __import__ 替换为安全版本
'__import__': self._safe_import,
# 注意:不包含 eval, exec, compile, open, input
}
return safe_builtins
def create_safe_exec_env(extra_libraries: Optional[List[str]] = None) -> dict:
"""创建安全的 exec 执行环境"""
importer = SecureImporter(extra_libraries)
return {
'__builtins__': importer.create_safe_builtins(),
'__name__': '__safe_module__',
}
if __name__ == "__main__":
code = """
import math
def test():
return math.sqrt(16)
"""
try:
exec(code, create_safe_exec_env())
print("✅ 标准库导入成功")
except Exception as e:
print(f"❌ 失败: {e}")core/operator_registry.py
"""共享算子注册中心(单例模式,全局唯一)"""
import hashlib
import json
from datetime import datetime
from typing import Dict, List, Optional, Any
from core.exceptions import OperatorExistsError
class OperatorRegistry:
"""
共享算子注册中心
特性:
- 单例模式,全局唯一
- 算子名称全局唯一,同名禁止覆盖
- 基于代码内容自动计算版本号
"""
_instance = None
_operators: Dict[str, Dict[str, Any]] = {} # key: op_name
def __new__(cls):
if cls._instance is None:
cls._instance = super().__new__(cls)
return cls._instance
def register(
self,
op_name: str,
code: str,
creator: str = "unknown",
description: str = "",
deps: List[str] = None
) -> str:
"""
注册一个新算子
Args:
op_name: 算子名称(全局唯一)
code: Python 函数代码(文本)
creator: 创建者标识
description: 算子功能描述
deps: 依赖库列表
Returns:
version: 版本号(基于内容哈希)
Raises:
OperatorExistsError: 同名算子已存在时抛出
"""
# 检查是否已存在(核心逻辑:禁止覆盖)
if op_name in self._operators:
existing = self._operators[op_name]
raise OperatorExistsError(
op_name=op_name,
creator=existing['creator'],
created_at=existing['created_at']
)
# 计算版本号
version = self._compute_version(code, deps)
# 存储算子
self._operators[op_name] = {
'name': op_name,
'code': code,
'deps': deps or [],
'creator': creator,
'description': description,
'created_at': datetime.now().isoformat(),
'version': version
}
return version
def get(self, op_name: str) -> Optional[Dict[str, Any]]:
"""获取算子信息"""
return self._operators.get(op_name)
def list_all(self) -> List[Dict[str, Any]]:
"""列出所有已注册算子"""
return list(self._operators.values())
def exists(self, op_name: str) -> bool:
"""检查算子是否存在"""
return op_name in self._operators
def _compute_version(self, code: str, deps: List[str]) -> str:
"""基于代码内容计算版本号"""
content = json.dumps({
'code': code,
'deps': sorted(deps or [])
}, sort_keys=True)
return hashlib.md5(content.encode()).hexdigest()
# 全局单例
registry = OperatorRegistry()core/dynamic_loader.py
import ast
import inspect
from typing import Dict, Optional, Any, List
import duckdb
from duckdb.sqltypes import VARCHAR, BIGINT, DOUBLE
from core.exceptions import OperatorNotFoundError, CodeCompileError
from core.operator_registry import registry
from core.security import create_safe_exec_env, SecurityError
# ================================================================
# AST 静态分析
# ================================================================
def _has_dangerous_reference(code: str) -> bool:
"""
检查代码中是否引用了危险内置函数或模块
包括:eval, exec, compile, open, input, __import__
"""
dangerous_names = {'eval', 'exec', 'compile', 'open', 'input', '__import__'}
try:
tree = ast.parse(code)
for node in ast.walk(tree):
# 检查函数调用:eval(...) 或 open(...)
if isinstance(node, ast.Call):
if isinstance(node.func, ast.Name) and node.func.id in dangerous_names:
return True
# 属性调用如 os.open, __builtins__.eval
if isinstance(node.func, ast.Attribute) and node.func.attr in dangerous_names:
return True
# 检查变量引用:eval(可能作为参数传递)
if isinstance(node, ast.Name) and node.id in dangerous_names:
# 如果它是函数调用的一部分,已经被上面的 Call 捕获
# 但为了安全,我们也标记单独的引用(如 x = eval)
return True
except SyntaxError:
# 如果代码有语法错误,交给 exec 处理
pass
return False
class DynamicUDFLoader:
"""
动态加载器:将用户注册的Python代码编译为DuckDB UDF
功能:
- 根据函数返回类型注解自动映射DuckDB类型
- 所有参数统一使用VARCHAR(DuckDB会隐式转换)
- 支持热加载和版本管理
"""
_loaded_versions: Dict[str, str] = {} # op_name -> version
@classmethod
def _compile_code(cls, code: str, op_name: str, extra_libraries: Optional[List[str]] = None) -> callable:
"""
在安全环境中编译用户代码
Args:
code: Python 函数代码文本
op_name: 函数名
extra_libraries: 额外允许的第三方库列表
Returns:
编译后的函数对象
"""
# 1. 静态检查:拦截危险函数引用
if _has_dangerous_reference(code):
raise CodeCompileError("检测到代码中包含危险函数引用 (eval, open, exec, compile, input, __import__)")
# 2. 执行安全环境
safe_env = create_safe_exec_env(extra_libraries)
try:
exec(code, safe_env)
except SecurityError as e:
raise CodeCompileError(f"安全限制: {e}")
except Exception as e:
raise CodeCompileError(f"代码编译失败: {e}")
if op_name not in safe_env:
raise CodeCompileError(f"代码中未找到函数: {op_name}")
func = safe_env[op_name]
if not callable(func):
raise CodeCompileError(f"{op_name} 不是可调用函数")
return func
@classmethod
def _map_python_type_to_duckdb(cls, py_type) -> Any:
"""
将Python类型映射为DuckDB类型
支持:
- int -> BIGINT
- float -> DOUBLE
- str -> VARCHAR
- 其他 -> VARCHAR(默认)
"""
if py_type is int:
return BIGINT
elif py_type is float:
return DOUBLE
elif py_type is str:
return VARCHAR
else:
return VARCHAR
@classmethod
def _get_function_signature(cls, func: callable):
"""
获取函数参数和返回类型
- 参数统一使用VARCHAR(DuckDB自动转换)
- 返回类型根据注解映射,无注解则默认为VARCHAR
"""
sig = inspect.signature(func)
num_params = len(sig.parameters)
param_types = [VARCHAR] * num_params
# 处理返回类型
return_annotation = sig.return_annotation
if return_annotation is inspect.Signature.empty:
return_type = VARCHAR
else:
return_type = cls._map_python_type_to_duckdb(return_annotation)
return param_types, return_type
@classmethod
def load_and_register(
cls,
conn: duckdb.DuckDBPyConnection,
op_name: str,
extra_libraries: Optional[List[str]] = None
) -> bool:
"""
加载并注册单个算子到DuckDB连接
Args:
conn: DuckDB连接对象
op_name: 算子名称
Returns:
True表示加载成功,False表示失败
"""
op_info = registry.get(op_name)
if not op_info:
raise OperatorNotFoundError(op_name)
current_ver = op_info['version']
if cls._loaded_versions.get(op_name) == current_ver:
return True
try:
func = cls._compile_code(op_info['code'], op_name, extra_libraries)
except Exception as e:
print(f"❌ 编译失败 {op_name}: {e}")
return False
param_types, return_type = cls._get_function_signature(func)
try:
conn.create_function(
name=op_name,
function=func,
parameters=param_types,
return_type=return_type,
type='native'
)
cls._loaded_versions[op_name] = current_ver
print(f"✅ 算子加载成功: {op_name} (v{current_ver[:8]})")
return True
except Exception as e:
print(f"❌ 注册到DuckDB失败 {op_name}: {e}")
return False
@classmethod
def load_all(cls, conn: duckdb.DuckDBPyConnection) -> int:
"""加载所有已注册算子"""
ops = registry.list_all()
count = 0
for op in ops:
if cls.load_and_register(conn, op['name']):
count += 1
return count
@classmethod
def reload_all(cls, conn: duckdb.DuckDBPyConnection) -> int:
"""强制重新加载所有算子(版本变更时使用)"""
cls._loaded_versions.clear()
return cls.load_all(conn)
def validate_operator_code(code: str, extra_libraries: Optional[List[str]] = None) -> tuple[bool, str]:
"""
独立检测算子代码是否合规(不注册到 DuckDB)
Args:
code: 算子代码
extra_libraries: 额外允许的第三方库列表
Returns:
(is_valid: bool, error_message: str)
"""
try:
# 1. AST 静态检测
if _has_dangerous_reference(code):
return False, "检测到代码中包含危险函数引用 (eval, open, exec, compile, input, __import__)"
# 2. 白名单检测(实际执行)
safe_env = create_safe_exec_env(extra_libraries)
try:
exec(code, safe_env)
except SecurityError as e:
return False, f"安全限制: {e}"
# 3. 检查是否至少定义了一个函数
has_function = any(
callable(v) and not v.__name__.startswith('__')
for v in safe_env.values()
)
if not has_function:
return False, "代码中未定义任何可调用函数"
return True, "合规"
except SyntaxError as e:
return False, f"语法错误: {e}"
except Exception as e:
return False, f"未知错误: {e}"
core/duckdb_engine.py
"""DuckDB数据治理引擎 - 核心类"""
from pathlib import Path
from typing import Optional, List, Dict, Any, Union
import duckdb
from core.operator_registry import registry
from core.dynamic_loader import DynamicUDFLoader
from core.exceptions import CodeCompileError
class DataGovernanceEngine:
"""
DuckDB数据治理引擎
功能:
1. 管理本地文件(CSV、JSON、Parquet)
2. 注册和管理共享算子(UDF)
3. 执行SQL查询(自动应用已注册的算子)
使用示例:
engine = DataGovernanceEngine(data_dir="./data")
engine.register_operator("clean_text", code, creator="admin")
result = engine.sql("SELECT clean_text(content) FROM raw_data")
"""
def __init__(
self,
data_dir: Optional[Union[str, Path]] = None,
auto_load_operators: bool = True
):
"""
初始化治理引擎
Args:
data_dir: 数据文件根目录
auto_load_operators: 是否自动加载所有已注册算子
"""
# 1. 设置数据目录
self.data_dir = None
if data_dir:
self.data_dir = Path(data_dir)
self.data_dir.mkdir(parents=True, exist_ok=True)
print(f"📁 数据目录: {self.data_dir.absolute()}")
# 2. 创建DuckDB连接(内存数据库)
self.conn = duckdb.connect(':memory:')
print("🔌 DuckDB内存连接已就绪")
# 3. 自动加载所有算子
if auto_load_operators:
count = DynamicUDFLoader.load_all(self.conn)
print(f"📦 已加载 {count} 个共享算子")
# 4. 内部状态
self._project_config = None
self._pipeline_views = []
def create_pipeline(self) -> "GovernancePipeline":
"""
创建治理流水线
Returns:
GovernancePipeline: 流水线实例
"""
from core.pipeline import GovernancePipeline
return GovernancePipeline(self)
# ==================== 算子管理API ====================
def register_operator(
self,
op_name: str,
code: str,
creator: str = "unknown",
description: str = "",
deps: List[str] = None,
extra_libraries: Optional[List[str]] = None, # ✅ 新增
) -> str:
"""
注册一个新算子(全局共享,同名禁止覆盖)
Args:
op_name: 算子名称(全局唯一)
code: Python函数代码(文本)
creator: 创建者标识
description: 算子功能描述
deps: 依赖库列表
extra_libraries: 额外允许的第三方库列表(如 ['pandas', 'numpy'])
Returns:
version: 版本号
Raises:
OperatorExistsError: 同名算子已存在时抛出
CodeCompileError: 代码编译失败时抛出
"""
# 1. 先编译验证(确保语法正确)
try:
# ✅ 传递 extra_libraries 给 _compile_code
DynamicUDFLoader._compile_code(code, op_name, extra_libraries)
except Exception as e:
raise CodeCompileError(f"算子代码编译失败: {e}")
# 2. 注册到中心
version = registry.register(op_name, code, creator, description, deps)
# 3. 加载到DuckDB(传递 extra_libraries)
success = DynamicUDFLoader.load_and_register(self.conn, op_name, extra_libraries)
if not success:
# 加载失败,回滚注册
if op_name in registry._operators:
del registry._operators[op_name]
raise RuntimeError("算子加载到DuckDB失败,已回滚注册")
print(f"📝 算子注册完成: {op_name} (v{version[:8]})")
return version
def list_operators(self) -> List[Dict[str, Any]]:
"""列出所有已注册算子"""
return registry.list_all()
def get_operator(self, op_name: str) -> Optional[Dict[str, Any]]:
"""获取指定算子信息"""
return registry.get(op_name)
def operator_exists(self, op_name: str) -> bool:
"""检查算子是否存在"""
return registry.exists(op_name)
def reload_operators(self) -> int:
"""重新加载所有算子"""
return DynamicUDFLoader.reload_all(self.conn)
# ==================== 文件操作API ====================
def read_csv(
self,
file_path: str,
view_name: Optional[str] = None,
**kwargs
) -> str:
"""读取CSV文件并注册为临时视图"""
full_path = self._resolve_path(file_path)
if not full_path.exists():
raise FileNotFoundError(f"文件不存在: {full_path}")
if view_name is None:
view_name = full_path.stem.replace('-', '_').replace('.', '_')
# 构建读取SQL
base_sql = f"CREATE OR REPLACE VIEW {view_name} AS SELECT * FROM read_csv('{full_path}'"
if kwargs:
params = [f"{k} = {repr(v)}" for k, v in kwargs.items()]
base_sql += ", " + ", ".join(params)
base_sql += ")"
self.conn.execute(base_sql)
print(f"📄 已加载CSV: {full_path} -> 视图 {view_name}")
return view_name
def read_parquet(self, file_path: str, view_name: Optional[str] = None) -> str:
"""读取Parquet文件并注册为临时视图"""
full_path = self._resolve_path(file_path)
if not full_path.exists():
raise FileNotFoundError(f"文件不存在: {full_path}")
if view_name is None:
view_name = full_path.stem.replace('-', '_').replace('.', '_')
self.conn.execute(
f"CREATE OR REPLACE VIEW {view_name} AS SELECT * FROM read_parquet('{full_path}')"
)
print(f"📄 已加载Parquet: {full_path} -> 视图 {view_name}")
return view_name
def register_table(self, table_name: str, data: Any) -> str:
"""注册任意数据为临时表"""
self.conn.register(table_name, data)
print(f"📊 已注册临时表: {table_name}")
return table_name
# ==================== SQL执行API ====================
def sql(self, query: str) -> duckdb.DuckDBPyRelation:
"""执行SQL查询,返回结果"""
return self.conn.sql(query)
def execute(self, query: str) -> None:
"""执行SQL(不返回结果)"""
self.conn.execute(query)
def table(self, table_name: str) -> duckdb.DuckDBPyRelation:
"""获取表引用"""
return self.conn.table(table_name)
def show_tables(self) -> List[str]:
"""显示所有已注册的表/视图"""
result = self.conn.sql("SHOW TABLES").fetchall()
return [row[0] for row in result]
# ==================== 辅助方法 ====================
def _resolve_path(self, file_path: str) -> Path:
"""解析文件路径"""
path = Path(file_path)
if path.is_absolute():
return path
if self.data_dir:
return self.data_dir / path
return path
def close(self):
"""关闭DuckDB连接"""
self.conn.close()
print("🔌 DuckDB连接已关闭")
def __enter__(self):
return self
def __exit__(self, exc_type, exc_val, exc_tb):
self.close()
core/mysql_loader.py
"""从MySQL加载治理项目配置"""
from typing import Dict, List, Optional
from sqlalchemy import create_engine, text
from sqlalchemy.orm import sessionmaker
class MySQLConfigLoader:
"""
从MySQL加载治理配置
表结构:
- governance_project: 项目信息
- governance_operator: 算子定义
- governance_step: 步骤定义
- governance_project_operator: 项目-算子关联
"""
def __init__(self, db_url: str):
"""
Args:
db_url: SQLAlchemy连接串,如 'mysql+pymysql://user:pass@host:3306/governance'
"""
self.engine = create_engine(
db_url,
pool_size=5,
max_overflow=10,
pool_pre_ping=True
)
self.Session = sessionmaker(bind=self.engine)
def get_project(self, project_name: str) -> Optional[Dict]:
"""获取项目信息"""
with self.Session() as session:
result = session.execute(
text("""
SELECT id, project_name, project_desc, data_dir, status
FROM governance_project
WHERE project_name = :name AND status = 1
"""),
{"name": project_name}
).mappings().first()
return dict(result) if result else None
def get_operators_for_project(self, project_name: str) -> List[Dict]:
"""获取项目关联的所有算子"""
with self.Session() as session:
sql = """
SELECT DISTINCT o.id, o.op_name, o.op_code, o.op_desc, o.creator
FROM governance_operator o
JOIN governance_project_operator po ON po.operator_id = o.id
JOIN governance_project p ON p.id = po.project_id
WHERE p.project_name = :name AND o.status = 1 AND p.status = 1
"""
results = session.execute(text(sql), {"name": project_name}).mappings().all()
return [dict(r) for r in results]
def get_steps(self, project_name: str) -> List[Dict]:
"""获取项目的所有步骤(按顺序)"""
with self.Session() as session:
sql = """
SELECT
s.id,
s.step_order,
s.step_name,
s.step_type,
s.sql_template,
s.output_view,
s.description
FROM governance_step s
JOIN governance_project p ON p.id = s.project_id
WHERE p.project_name = :name AND p.status = 1
ORDER BY s.step_order ASC
"""
results = session.execute(text(sql), {"name": project_name}).mappings().all()
return [dict(r) for r in results]
def get_full_project_config(self, project_name: str) -> Dict:
"""获取完整项目配置(项目信息 + 算子 + 步骤)"""
project = self.get_project(project_name)
if not project:
raise ValueError(f"项目不存在或已禁用: {project_name}")
operators = self.get_operators_for_project(project_name)
steps = self.get_steps(project_name)
return {
"project": project,
"operators": operators,
"steps": steps
}
core/pipeline.py
"""治理流水线 - 按步骤执行SQL,视图透传"""
from typing import Dict, List, Optional, Any
import duckdb
from core.duckdb_engine import DataGovernanceEngine
class GovernancePipeline:
"""
治理流水线:按步骤执行SQL,前一步结果作为视图传递给下一步
支持变量:
- {view}: 上一步的输出视图名
- {{file_path}}: 文件路径(由运行时传入)
使用示例:
pipeline = GovernancePipeline(engine)
pipeline.add_step("load", "CREATE OR REPLACE VIEW raw AS SELECT * FROM read_csv('{{file_path}}')")
pipeline.add_step("clean", "SELECT * FROM {view} WHERE ...")
result = pipeline.run()
"""
def __init__(self, engine: DataGovernanceEngine):
self.engine = engine
self.steps: List[Dict[str, Any]] = []
self._views: List[str] = []
self._file_path: Optional[str] = None
def set_file_path(self, file_path: str) -> "GovernancePipeline":
"""设置数据文件路径(自动解析为绝对路径)"""
# 使用引擎的 _resolve_path 方法将相对路径转为绝对路径
resolved = self.engine._resolve_path(file_path)
self._file_path = str(resolved)
print(f"📄 数据文件路径解析: {file_path} -> {self._file_path}")
return self
def add_step(
self,
name: str,
sql: str,
output_view: Optional[str] = None,
description: str = ""
) -> "GovernancePipeline":
"""
添加一个治理步骤
Args:
name: 步骤名称(用于日志)
sql: SQL语句(支持{view}和{{file_path}}变量)
output_view: 输出视图名称(如不指定则不保存)
description: 步骤描述
"""
self.steps.append({
"name": name,
"sql": sql,
"output_view": output_view,
"description": description
})
return self
def _render_sql(self, sql: str, context: Dict[str, str]) -> str:
"""
渲染SQL模板
- {view} -> 上一步的输出视图
- {{file_path}} -> 设置的文件路径
"""
rendered = sql
# 替换 {view} 为上一步视图
if "{view}" in rendered and context.get("last_view"):
rendered = rendered.replace("{view}", context["last_view"])
# 替换 {{file_path}} 为文件路径
if "{{file_path}}" in rendered and self._file_path:
rendered = rendered.replace("{{file_path}}", self._file_path)
return rendered
def run(self, show_progress: bool = True) -> Optional[duckdb.DuckDBPyRelation]:
"""
执行流水线
Returns:
最后一步的结果(DuckDBPyRelation),若无结果则返回None
"""
context = {
"last_view": None,
"all_views": []
}
result = None
total = len(self.steps)
for idx, step in enumerate(self.steps, 1):
name = step["name"]
sql = step["sql"]
output_view = step.get("output_view")
desc = step.get("description", "")
if show_progress:
print(f" [{idx}/{total}] {name} - {desc or '执行中...'}")
# 检查是否还有未替换的 {{file_path}}
if "{{file_path}}" in sql and self._file_path is None:
raise ValueError(f"步骤 '{name}' 需要设置 file_path,请调用 set_file_path()")
# 渲染SQL
rendered_sql = self._render_sql(sql, context)
try:
if output_view:
# 【修复】检查SQL是否已经包含CREATE语句
# 如果已包含,直接执行;否则包装CREATE OR REPLACE VIEW
sql_upper = rendered_sql.strip().upper()
if sql_upper.startswith("CREATE") or sql_upper.startswith("INSERT"):
full_sql = rendered_sql
else:
full_sql = f"CREATE OR REPLACE VIEW {output_view} AS {rendered_sql}"
self.engine.execute(full_sql)
context["last_view"] = output_view
context["all_views"].append(output_view)
self._views.append(output_view)
result = self.engine.sql(f"SELECT * FROM {output_view} LIMIT 5")
else:
# 直接执行,不保存视图(通常用于最终输出)
result = self.engine.sql(rendered_sql)
except Exception as e:
print(f" ❌ 步骤 '{name}' 执行失败: {e}")
print(f" SQL: {rendered_sql[:300]}...")
raise
if context["last_view"]:
return self.engine.sql(f"SELECT * FROM {context['last_view']}")
return result
def get_views(self) -> List[str]:
"""获取所有创建的视图名称"""
return self._views
def print_summary(self):
"""打印流水线摘要"""
print("\n📊 流水线执行摘要:")
print(f" 总步骤数: {len(self.steps)}")
print(f" 创建视图: {', '.join(self._views) if self._views else '无'}")
core/duckdb_engine_ext.py
"""扩展DataGovernanceEngine,增加MySQL配置加载能力"""
from pathlib import Path
from typing import Optional
from core.duckdb_engine import DataGovernanceEngine
from core.mysql_loader import MySQLConfigLoader
from core.dynamic_loader import DynamicUDFLoader
from core.operator_registry import registry
from core.pipeline import GovernancePipeline
class ExtendedGovernanceEngine(DataGovernanceEngine):
"""
扩展治理引擎:增加MySQL配置加载和流水线执行能力
"""
def load_project_from_mysql(
self,
db_url: str,
project_name: str,
data_dir: Optional[str] = None
) -> "ExtendedGovernanceEngine":
"""
从MySQL加载项目配置(算子 + 步骤)
Args:
db_url: MySQL连接串
project_name: 项目名称
data_dir: 数据目录(覆盖MySQL配置)
"""
loader = MySQLConfigLoader(db_url)
config = loader.get_full_project_config(project_name)
# 1. 设置数据目录
if data_dir:
self.data_dir = Path(data_dir)
elif config["project"].get("data_dir"):
self.data_dir = Path(config["project"]["data_dir"])
if self.data_dir:
self.data_dir.mkdir(parents=True, exist_ok=True)
print(f"📁 数据目录: {self.data_dir.absolute()}")
# 2. 注册算子
print(f"\n【加载算子】项目: {project_name}")
for op in config["operators"]:
try:
# 先编译验证
DynamicUDFLoader._compile_code(op["op_code"], op["op_name"])
# 注册到注册中心
registry.register(
op_name=op["op_name"],
code=op["op_code"],
creator=op.get("creator", "system"),
description=op.get("op_desc", "")
)
# 加载到DuckDB
DynamicUDFLoader.load_and_register(self.conn, op["op_name"])
print(f" ✅ 加载算子: {op['op_name']}")
except Exception as e:
print(f" ⚠️ 算子 {op['op_name']} 加载失败: {e}")
# 3. 保存步骤配置
self._project_config = {
"name": project_name,
"steps": config["steps"],
"project_info": config["project"]
}
return self
def create_pipeline_from_mysql(
self,
db_url: str,
project_name: str,
file_path: Optional[str] = None,
data_dir: Optional[str] = None
) -> GovernancePipeline:
"""
从MySQL创建治理流水线
Args:
db_url: MySQL连接串
project_name: 项目名称
file_path: 数据文件路径(替换{{file_path}})
data_dir: 数据目录
Returns:
配置好的GovernancePipeline实例
"""
# 加载项目配置
self.load_project_from_mysql(db_url, project_name, data_dir)
# 创建流水线
pipeline = GovernancePipeline(self)
if file_path:
pipeline.set_file_path(file_path)
# 添加步骤
for step in self._project_config["steps"]:
sql = step["sql_template"]
# 检查是否包含需要替换的变量
if "{{file_path}}" in sql and not file_path:
print(f" ⚠️ 步骤 '{step['step_name']}' 需要 file_path 参数,请传入")
pipeline.add_step(
name=step["step_name"],
sql=sql,
output_view=step.get("output_view"),
description=step.get("description", "")
)
return pipeline
def run_project_from_mysql(
self,
db_url: str,
project_name: str,
file_path: str,
data_dir: Optional[str] = None
):
"""
一键执行MySQL配置的治理项目
Args:
db_url: MySQL连接串
project_name: 项目名称
file_path: 数据文件路径(必填)
data_dir: 数据目录(可选)
Returns:
治理结果
"""
pipeline = self.create_pipeline_from_mysql(db_url, project_name, file_path, data_dir)
print(f"\n【执行治理流水线】项目: {project_name}")
print(f" 数据文件: {file_path}")
result = pipeline.run()
# 打印摘要
pipeline.print_summary()
return result
core/init.py
"""数据治理平台核心模块"""
from core.exceptions import (
OperatorExistsError,
OperatorNotFoundError,
CodeCompileError,
StepExecutionError
)
from core.operator_registry import OperatorRegistry, registry
from core.dynamic_loader import DynamicUDFLoader
from core.duckdb_engine import DataGovernanceEngine
from core.duckdb_engine_ext import ExtendedGovernanceEngine
from core.pipeline import GovernancePipeline
from core.mysql_loader import MySQLConfigLoader
__all__ = [
"OperatorExistsError",
"OperatorNotFoundError",
"CodeCompileError",
"StepExecutionError",
"OperatorRegistry",
"registry",
"DynamicUDFLoader",
"DataGovernanceEngine",
"ExtendedGovernanceEngine",
"GovernancePipeline",
"MySQLConfigLoader"
]
示例:
examples/generate_data.py
"""生成示例CSV数据(包含脏数据)"""
import csv
import random
from datetime import datetime, timedelta
# 固定随机种子
random.seed(42)
def generate_sample_csv(file_path: str, rows: int = 100):
"""生成包含脏数据的示例CSV"""
sources = ['web', 'pdf', 'excel', 'api', 'manual']
categories = ['finance', 'medical', 'tech', 'legal', 'education']
# 预置重复文本
duplicates = [
"这是完全重复的文本A,用于测试去重",
"这是完全重复的文本B,用于测试去重",
]
with open(file_path, 'w', newline='', encoding='utf-8-sig') as f:
writer = csv.writer(f)
writer.writerow(['id', 'source', 'category', 'content', 'create_date', 'score'])
for i in range(1, rows + 1):
r = random.random()
if r < 0.08:
content = None
score = None
elif r < 0.15:
content = "\t\t\t" if random.random() > 0.5 else ""
score = 0.0
elif r < 0.22:
content = random.choice(duplicates)
score = round(random.uniform(0.5, 0.9), 2)
else:
# 正常文本
template = random.choice([
"根据{year}年财报显示,{company}的营收达到{amount}亿元,同比增长{growth}%。",
"患者{name},{age}岁,主诉{symptom},初步诊断为{diagnosis}。",
"技术文档:{system} v{version} 修复了{bug},提升了{performance}%。",
"合同编号{code},甲方{party_a}与乙方{party_b}就{project}达成协议。"
])
content = template.format(
year=random.randint(2020, 2026),
company=random.choice(['华为', '腾讯', '阿里', '字节', '百度']),
amount=round(random.uniform(10, 500), 1),
growth=round(random.uniform(5, 35), 1),
name=random.choice(['张三', '李四', '王五', '赵六']),
age=random.randint(20, 70),
symptom=random.choice(['头痛', '发热', '咳嗽', '胸痛', '无症状']),
diagnosis=random.choice(['感冒', '流感', '肺炎', '高血压', '健康']),
system=random.choice(['Kubernetes', 'MySQL', 'Redis', 'Nginx']),
version=f"{random.randint(1, 3)}.{random.randint(0, 9)}.{random.randint(0, 9)}",
bug=random.choice(['内存泄漏', '死锁', '缓存穿透', '安全漏洞']),
performance=random.randint(10, 60),
code=f"CT-{random.randint(1000, 9999)}",
party_a=random.choice(['中国移动', '中国电信', '国家电网']),
party_b=random.choice(['华为', '中兴', '浪潮']),
project=random.choice(['智慧城市', '数据中心', '5G基站'])
)
if random.random() > 0.8:
content = content.replace("。", "。\n")
score = round(random.uniform(0.3, 0.99), 2)
if content in (None, "", "\t\t\t"):
score = None
writer.writerow([
i,
random.choice(sources),
random.choice(categories),
content,
(datetime.now() - timedelta(days=random.randint(0, 90))).strftime('%Y-%m-%d'),
score
])
print(f"✅ 示例数据已生成: {file_path} ({rows} 条记录)")
if __name__ == "__main__":
generate_sample_csv("./data/sample_data.csv", 100)
examples/demo_pipeline.py
"""流水线模式演示"""
import os
import sys
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from core.duckdb_engine import DataGovernanceEngine
from core.exceptions import OperatorExistsError
def main():
print("=" * 60)
print("🚀 治理流水线演示")
print("=" * 60)
# 1. 初始化引擎
engine = DataGovernanceEngine(data_dir="./data")
# 2. 注册算子
try:
engine.register_operator(
"clean_text",
"def clean_text(text: str) -> str:\n if text is None:\n return ''\n return ' '.join(text.split())",
creator="admin"
)
except OperatorExistsError:
pass
try:
engine.register_operator(
"text_length",
"def text_length(text: str) -> int:\n if text is None:\n return 0\n return len(text.replace(' ', ''))",
creator="admin"
)
except OperatorExistsError:
pass
# 3. 创建流水线
pipeline = engine.create_pipeline()
pipeline.set_file_path("*.csv")
# 4. 添加步骤
pipeline.add_step(
"加载数据",
"CREATE OR REPLACE VIEW raw_data AS SELECT * FROM read_csv('{{file_path}}', header=True)",
output_view="raw_data",
description="加载CSV"
)
pipeline.add_step(
"清洗文本",
"CREATE OR REPLACE VIEW cleaned AS SELECT *, clean_text(content) AS cleaned_content, text_length(content) AS content_len FROM {view}",
output_view="cleaned",
description="清洗文本"
)
pipeline.add_step(
"过滤短文本",
"CREATE OR REPLACE VIEW filtered AS SELECT * FROM {view} WHERE content_len >= 10",
output_view="filtered",
description="过滤短文本"
)
pipeline.add_step(
"去重",
"CREATE OR REPLACE VIEW deduped AS SELECT * FROM {view} QUALIFY ROW_NUMBER() OVER (PARTITION BY cleaned_content ORDER BY score DESC NULLS LAST) = 1",
output_view="deduped",
description="去重"
)
pipeline.add_step(
"生成报告",
"SELECT source, category, COUNT(*) as total, AVG(content_len) as avg_len FROM {view} GROUP BY source, category ORDER BY source, total DESC",
output_view=None,
description="生成报告"
)
# 5. 执行流水线
print("\n执行流水线...")
result = pipeline.run()
print("\n治理结果:")
print(result.df().to_string(index=False))
print(f"\n创建的视图: {pipeline.get_views()}")
engine.close()
if __name__ == "__main__":
main()
init_mysql.sql
-- ================================================================
-- 数据库: governance_platform
-- 描述: 治理平台配置存储
-- ================================================================
-- 创建数据库
CREATE
DATABASE IF NOT EXISTS `governance_platform`
DEFAULT CHARACTER SET utf8mb4
COLLATE utf8mb4_unicode_ci;
USE
`governance_platform`;
-- ----------------------------------------------------------------
-- 1. 治理项目表
-- ----------------------------------------------------------------
DROP TABLE IF EXISTS `governance_project`;
CREATE TABLE `governance_project`
(
`id` INT PRIMARY KEY AUTO_INCREMENT COMMENT '项目ID',
`project_name` VARCHAR(100) NOT NULL COMMENT '项目名称(唯一标识)',
`project_desc` VARCHAR(500) DEFAULT '' COMMENT '项目描述',
`data_dir` VARCHAR(500) DEFAULT '' COMMENT '数据文件目录',
`status` TINYINT DEFAULT 1 COMMENT '状态: 0-禁用, 1-启用',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
UNIQUE KEY `uk_project_name` (`project_name`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='治理项目表';
-- ----------------------------------------------------------------
-- 2. 算子定义表(共享算子池)
-- ----------------------------------------------------------------
DROP TABLE IF EXISTS `governance_operator`;
CREATE TABLE `governance_operator`
(
`id` INT PRIMARY KEY AUTO_INCREMENT COMMENT '算子ID',
`op_name` VARCHAR(100) NOT NULL COMMENT '算子名称(SQL中调用)',
`op_code` LONGTEXT NOT NULL COMMENT 'Python函数代码',
`op_desc` VARCHAR(500) DEFAULT '' COMMENT '算子描述',
`creator` VARCHAR(100) DEFAULT 'admin' COMMENT '创建者',
`status` TINYINT DEFAULT 1 COMMENT '状态: 0-禁用, 1-启用',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
UNIQUE KEY `uk_op_name` (`op_name`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='算子定义表';
-- ----------------------------------------------------------------
-- 3. 治理步骤表
-- ----------------------------------------------------------------
DROP TABLE IF EXISTS `governance_step`;
CREATE TABLE `governance_step`
(
`id` INT PRIMARY KEY AUTO_INCREMENT COMMENT '步骤ID',
`project_id` INT NOT NULL COMMENT '所属项目ID',
`step_order` INT NOT NULL COMMENT '执行顺序(从1开始)',
`step_name` VARCHAR(100) NOT NULL COMMENT '步骤名称',
`step_type` ENUM('load', 'sql') DEFAULT 'sql' COMMENT '步骤类型: load-加载, sql-SQL执行',
`sql_template` LONGTEXT NOT NULL COMMENT 'SQL模板(支持{view}和{{file_path}}变量)',
`output_view` VARCHAR(100) DEFAULT NULL COMMENT '输出视图名',
`description` VARCHAR(500) DEFAULT '' COMMENT '步骤描述',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
FOREIGN KEY (`project_id`) REFERENCES `governance_project` (`id`) ON DELETE CASCADE,
UNIQUE KEY `uk_project_step` (`project_id`, `step_order`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='治理步骤表';
-- ----------------------------------------------------------------
-- 4. 项目-算子关联表
-- ----------------------------------------------------------------
DROP TABLE IF EXISTS `governance_project_operator`;
CREATE TABLE `governance_project_operator`
(
`id` INT PRIMARY KEY AUTO_INCREMENT COMMENT '关联ID',
`project_id` INT NOT NULL COMMENT '项目ID',
`operator_id` INT NOT NULL COMMENT '算子ID',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
FOREIGN KEY (`project_id`) REFERENCES `governance_project` (`id`) ON DELETE CASCADE,
FOREIGN KEY (`operator_id`) REFERENCES `governance_operator` (`id`) ON DELETE CASCADE,
UNIQUE KEY `uk_project_op` (`project_id`, `operator_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='项目算子关联表';
-- ================================================================
-- 初始化示例数据
-- ================================================================
-- 插入示例算子
INSERT INTO `governance_operator` (`op_name`, `op_code`, `op_desc`, `creator`)
VALUES ('clean_text',
'def clean_text(text: str) -> str:\n if text is None:\n return ""\n return " ".join(text.split())',
'去除首尾空白,合并多个空格为单个空格',
'admin'),
('text_length',
'def text_length(text: str) -> int:\n if text is None:\n return 0\n return len(text.replace(" ", "").replace("\\n", ""))',
'计算有效文本长度(不含空格和换行)',
'admin'),
('extract_numbers',
'import re\ndef extract_numbers(text: str) -> str:\n if text is None:\n return ""\n nums = re.findall(r"-?\\d+\\.?\\d*", text)\n return ",".join(nums)',
'从文本中提取所有数字,用逗号分隔',
'admin');
-- 插入示例项目
INSERT INTO `governance_project` (`project_name`, `project_desc`, `data_dir`)
VALUES ('数据质量治理', '清洗CSV数据、去重、生成质量报告', './data');
-- 插入项目-算子关联
INSERT INTO `governance_project_operator` (`project_id`, `operator_id`)
SELECT 1, id
FROM `governance_operator`
WHERE op_name IN ('clean_text', 'text_length', 'extract_numbers');
-- 插入示例步骤
INSERT INTO `governance_step` (`project_id`, `step_order`, `step_name`, `step_type`, `sql_template`, `output_view`,
`description`)
VALUES (1, 1, '加载数据', 'load',
'CREATE OR REPLACE VIEW raw_data AS SELECT * FROM read_csv(''{{file_path}}'', header=True)',
'raw_data',
'加载CSV文件'),
(1, 2, '清洗文本', 'sql',
'CREATE OR REPLACE VIEW cleaned AS SELECT *, clean_text(content) AS cleaned_content, text_length(content) AS content_len FROM {view}',
'cleaned',
'应用clean_text算子清洗文本'),
(1, 3, '过滤短文本', 'sql',
'CREATE OR REPLACE VIEW filtered AS SELECT * FROM {view} WHERE content_len >= 10 AND content IS NOT NULL',
'filtered',
'过滤文本长度小于10的记录'),
(1, 4, '去重', 'sql',
'CREATE OR REPLACE VIEW deduped AS SELECT * FROM {view} QUALIFY ROW_NUMBER() OVER (PARTITION BY cleaned_content ORDER BY score DESC NULLS LAST) = 1',
'deduped',
'按cleaned_content去重,保留质量分最高'),
(1, 5, '生成质量报告', 'sql',
'SELECT source, category, COUNT(*) AS total, AVG(content_len) AS avg_len, SUM(CASE WHEN score >= 0.7 THEN 1 ELSE 0 END) AS high_quality FROM {view} GROUP BY source, category ORDER BY source, total DESC',
NULL,
'按来源和分类生成聚合报告');examples/demo_mysql.py
"""MySQL配置驱动治理演示"""
import os
import sys
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from core.duckdb_engine_ext import ExtendedGovernanceEngine
from core.mysql_loader import MySQLConfigLoader
def main():
print("=" * 60)
print("🚀 MySQL配置驱动治理演示")
print("=" * 60)
# MySQL连接配置
DB_URL = "mysql+pymysql://root:root@localhost:3307/governance_platform"
# 1. 初始化扩展引擎
engine = ExtendedGovernanceEngine()
# 2. 一键执行项目
result = engine.run_project_from_mysql(
db_url=DB_URL,
project_name="数据质量治理",
file_path="*.csv",
data_dir="./data"
)
# 3. 查看结果
print("\n【治理结果】")
if result is not None:
print(result.df().to_string(index=False))
else:
print("无结果返回")
# 4. 查看已注册算子
print("\n【当前已注册算子】")
for op in engine.list_operators():
print(f" 📌 {op['name']} (v{op['version'][:8]})")
engine.close()
def demo_manual_builder():
"""手动构建流水线(适合调试)"""
print("\n" + "=" * 60)
print("🔧 手动构建流水线(调试模式)")
print("=" * 60)
DB_URL = "mysql+pymysql://root:root@localhost:3307/governance_platform"
engine = ExtendedGovernanceEngine()
# 加载配置
loader = MySQLConfigLoader(DB_URL)
config = loader.get_full_project_config("数据质量治理")
print(f"\n📋 项目: {config['project']['project_name']}")
print(f" 算子数: {len(config['operators'])}")
print(f" 步骤数: {len(config['steps'])}")
engine.close()
if __name__ == "__main__":
# 先确保数据存在
# from generate_data import generate_sample_csv
# generate_sample_csv("../data/sample_data.csv", 100)
main()
# demo_manual_builder()
examples/demo_security.py
"""
算子安全机制测试 Demo
核心测试目标:
1. 标准库(re, math, datetime) → 应该成功注册
2. pandas + numpy + extra_libraries → 应该成功注册
3. pandas 不传递 extra_libraries → 应该被拦截
4. 危险库(os, subprocess, socket, __import__, eval, open, os.path) → 应该被拦截
"""
import os
import sys
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from core.duckdb_engine import DataGovernanceEngine
from core.exceptions import CodeCompileError
def print_separator(title: str = "", char: str = "=", length: int = 60):
if title:
print(f"\n{char * 5} {title} {char * (length - len(title) - 7)}")
else:
print(char * length)
# ---------- 测试1:标准库 ----------
def test_legal_operators():
print_separator("测试1:合法算子(标准库)")
engine = DataGovernanceEngine(auto_load_operators=False)
code1 = """
import re
import math
def process_text(text: str) -> str:
if text is None:
return ''
cleaned = re.sub(r'[^\\w\\s]', '', text)
log_len = math.log(len(cleaned) + 1)
return f"{cleaned[:50]}... (length_log: {log_len:.2f})"
"""
code2 = """
from datetime import datetime
def get_current_time() -> str:
return datetime.now().isoformat()
"""
try:
engine.register_operator("process_text", code1, creator="admin")
engine.register_operator("get_current_time", code2, creator="admin")
print("✅ 合法算子(标准库)注册成功!")
except CodeCompileError as e:
print(f"❌ 意外失败: {e}")
finally:
engine.close()
# ---------- 测试2:pandas + numpy + extra_libraries ----------
def test_pandas_with_extra_libraries():
print_separator("测试2:pandas + numpy(传递 extra_libraries)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
import pandas as pd
import numpy as np
def series_sum(data: list) -> float:
if not data:
return 0.0
return pd.Series(data).sum()
"""
try:
engine.register_operator(
"series_sum",
code,
creator="admin",
extra_libraries=['pandas', 'numpy']
)
print("✅ pandas + numpy 注册成功(extra_libraries 生效)")
except CodeCompileError as e:
print(f"❌ 意外失败: {e}")
finally:
engine.close()
# ---------- 测试3:pandas 但不传 extra_libraries ----------
def test_pandas_without_extra_libraries():
print_separator("测试3:pandas 不传 extra_libraries(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
import pandas as pd
def test_func() -> str:
return "test"
"""
try:
engine.register_operator("test_func", code, creator="admin")
print("❌ 危险!pandas 被错误允许了")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试4:os ----------
def test_illegal_os():
print_separator("测试4:导入 os(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
import os
def delete_file(path: str) -> str:
os.remove(path)
return "deleted"
"""
try:
engine.register_operator("delete_file", code, creator="admin")
print("❌ os 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试5:subprocess ----------
def test_illegal_subprocess():
print_separator("测试5:导入 subprocess(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
import subprocess
def run(cmd: str) -> str:
return subprocess.run(cmd, shell=True).stdout
"""
try:
engine.register_operator("run", code, creator="admin")
print("❌ subprocess 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试6:socket ----------
def test_illegal_socket():
print_separator("测试6:导入 socket(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
import socket
def hostname() -> str:
return socket.gethostname()
"""
try:
engine.register_operator("hostname", code, creator="admin")
print("❌ socket 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试7:__import__ ----------
def test_illegal_import():
print_separator("测试7:动态导入 __import__(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
def dangerous_import(name: str):
return __import__(name)
"""
try:
engine.register_operator("dangerous_import", code, creator="admin")
print("❌ __import__ 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试8:eval ----------
def test_illegal_eval():
print_separator("测试8:使用 eval(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
def dangerous_eval(code: str):
return eval(code)
"""
try:
engine.register_operator("dangerous_eval", code, creator="admin")
print("❌ eval 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试9:open ----------
def test_illegal_open():
print_separator("测试9:使用 open(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
def read_file(path: str) -> str:
with open(path, 'r') as f:
return f.read()
"""
try:
engine.register_operator("read_file", code, creator="admin")
print("❌ open 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 测试10:os.path ----------
def test_illegal_os_path():
print_separator("测试10:导入 os.path(应被拦截)")
engine = DataGovernanceEngine(auto_load_operators=False)
code = """
import os.path
def check(path: str) -> bool:
return os.path.exists(path)
"""
try:
engine.register_operator("check", code, creator="admin")
print("❌ os.path 被错误允许")
except CodeCompileError as e:
print(f"✅ 成功拦截: {e}")
finally:
engine.close()
# ---------- 主函数 ----------
def main():
print("=" * 60)
print("🔒 算子安全机制测试(专注拦截验证)")
print("=" * 60)
print("\n📋 测试项:")
print(" 1. 标准库 → 应成功")
print(" 2. pandas + extra_libraries → 应成功")
print(" 3. pandas 不传 extra_libraries → 应拦截")
print(" 4-10. 危险模块/函数 → 应拦截")
print("=" * 60)
test_legal_operators()
test_pandas_with_extra_libraries()
test_pandas_without_extra_libraries()
test_illegal_os()
test_illegal_subprocess()
test_illegal_socket()
test_illegal_import()
test_illegal_eval()
test_illegal_open()
test_illegal_os_path()
print("\n" + "=" * 60)
print("✅ 所有安全测试完成!")
print("=" * 60)
if __name__ == "__main__":
main()