数据治理

构建DuckDB 数据治理平台:算子引擎、安全机制与流水线设计

一、写在前面

在数据治理的日常工作中,我们经常需要处理各种非结构化数据——从 PDF 中提取的文本、爬虫抓取的网页内容、API 返回的 JSON 数据。这些数据往往包含空值、重复、格式混乱等问题,需要一个高效、灵活的数据清洗方案。

传统做法是写 Python 脚本,用 Pandas 进行清洗。但随着数据量增长(几十 GB 甚至更大),Pandas 的内存瓶颈开始显现——数据必须全部加载到内存中,稍有不慎就会 OOM。

于是我将目光投向了 DuckDB:一个嵌入式、列式存储、支持 SQL 的 OLAP 引擎。它可以直接读取 CSV、Parquet 文件,支持流式处理,数据量远超内存也能跑。

基于 DuckDB,我构建了一个完整的 数据治理平台,核心功能包括:

  • 算子引擎:用户可注册自定义 Python UDF,在 SQL 中直接调用
  • 治理流水线:按步骤执行 SQL,视图透传,支持链式调用
  • 安全机制:限制用户代码只能导入白名单模块,拦截危险函数
  • 配置驱动:MySQL 存储项目配置,支持前端在线编排

本文将分享这个项目的设计思路与核心实现。

二、为什么选择 DuckDB?

对比项PandasDuckDB
内存使用全量加载到内存流式处理,可处理 > 内存的数据
数据源CSV、ExcelCSV、JSON、Parquet、Excel、数据库
操作方式Python APISQL(支持标准 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 函数。我们将其封装为 算子引擎

核心设计

  1. 算子注册中心:单例模式,全局唯一,同名禁止覆盖
  2. 动态加载:从 MySQL 读取算子代码,通过 exec 动态编译并注册到 DuckDB
  3. 版本管理:基于代码内容自动生成版本号
# 使用示例
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(抽象语法树),检测是否包含危险函数引用(evalopenexec__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__,只允许导入白名单中的模块。白名单分为两类:

  • 标准库(安全子集)mathrejsondatetimepandas 等
  • 额外允许的第三方库:用户通过 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) → 去重(保留质量分最高) → 聚合报告

治理前后对比

指标治理前治理后
总行数10081
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
SQLAlchemyMySQL 连接池管理≥ 2.0.0
PyMySQLMySQL 驱动≥ 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()