ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

ezdata:TB级数据流水线的低代码操作系统

ezdata:TB级数据流水线的低代码操作系统 简介本资源是一个面向数据工程师与全栈开发者的开源数据平台完整实现聚焦TB级数据处理、多源统一建模与智能分析场景解决企业级数据中台建设中常见的异构数据接入、DAG任务调度、低代码集成及LLM增强问答等核心问题。压缩包共2000个文件主体为629个TypeScript含Vue3组件与状态管理逻辑、759个Vue单文件组件构建响应式前端界面、263个Python后端模块涵盖分布式Pandas引擎封装、任务调度器、LLM接口适配与数据模型抽象层辅以YAML配置、SVG图标及Dockerfile等工程化文件整体仅9.32MB结构紧凑且模块边界清晰。已有38人下载学习适合具备PythonVue基础、正实践数据平台搭建或研究低代码AI融合架构的中高级开发者。读者可直接运行完整前后端系统获取标准化多数据源接入流程、可复用的DAG编排SDK、统一Schema映射方案及LLM驱动的数据语义解析示例代码。1. 这不是又一个“前后端分离 demo”ezdata 系统本质是把 TB 级数据工程流水线塞进低代码界面让业务同学能拖拽调度 Pandas DAG、用自然语言查清洗结果、在同一个模型层里混接 MySQL/CSV/API/Parquet——它解决的不是“怎么写接口”而是“怎么让数据分析师不求人改 SQL 就能上线新指标”你见过凌晨三点还在改 Airflow DAG YAML 的数据工程师吗见过业务方拿着 Excel 表找开发“这个字段能不能加个去重逻辑”而你得切 Git 分支、写 Pandas UDF、提 PR、等 CI、发版、回滚……最后发现只是少了个.drop_duplicates()ezdata 不是炫技的全栈玩具。它是一套可落地的数据生产力操作系统后端用 Python非 Flask 快速原型而是 FastAPI 异步任务队列 分布式 Pandas 执行器前端用 Vue3非单页管理后台而是带 DAG 编排画布、LLM 问答输入框、多源连接器面板的低代码 IDE。核心不在“用了 Vue3”而在“用 Vue3 实现了 LLM 驱动的 Schema 探索”不在“支持多数据源”而在“所有数据源统一映射到ezdata.Table抽象连 CSV 的 header 和 PostgreSQL 的information_schema都被编译成同一份元数据 JSON”更关键的是——它的“低代码”不是拖拽生成 CRUD 页面而是拖拽定义数据血缘、自然语言触发清洗任务、点击导出为可复用的 Pandas Pipeline 模块。适合三类人需要快速验证数据逻辑的 BI 工程师、想绕过开发直接跑分析脚本的数据分析师、以及正被重复性 ETL 压垮的中小团队后端。如果你的痛点是“需求来了等开发排期两周”那 ezdata 的价值就不是技术选型而是时间成本重构。2. 后端FastAPI 分布式 Pandas 执行器不是“Python 写个 API”而是构建可水平伸缩的数据计算底座2.1 为什么选 FastAPI 而不是 Django 或 Flask——类型驱动的元数据契约与异步任务流天然契合ezdata 后端的核心约束不是“要快”而是“要可推导”。当用户在前端拖拽一个“过滤订单金额 1000”的节点时系统必须在提交前就校验该字段是否存在类型是否支持比较上游表是否已加载这要求 API 层具备强类型契约能力。FastAPI 的 Pydantic v2 模型天然支持嵌套结构、条件校验、JSON Schema 导出——我们定义TaskNodeBase为所有算子基类# schemas/task.py from pydantic import BaseModel, Field, validator from typing import Optional, Dict, Any, List class TaskNodeBase(BaseModel): id: str Field(..., description节点唯一ID前端生成) type: str Field(..., description算子类型filter, join, groupby, llm_query...) upstream: List[str] Field(default_factorylist, description上游节点ID列表) validator(type) def validate_type(cls, v): if v not in [filter, join, groupby, llm_query, pandas_udf]: raise ValueError(f不支持的算子类型: {v}) return v class FilterNode(TaskNodeBase): type: str filter column: str Field(..., description待过滤字段名) operator: str Field(..., description操作符eq, gt, lt, contains...) value: Any Field(..., description过滤值类型由column schema推导)提示value: Any看似危险但实际校验发生在FilterNode.parse_obj()时——Pydantic 会根据column在上游表 Schema 中的 dtype如int64,string自动 cast 并报错。这比 Flask manual JSON 解析早 3 步拦截错误。更重要的是异步支持。DAG 调度需并发执行多个 Pandas 任务如同时读取 MySQL 和 S3 Parquet且不能阻塞主线程。FastAPI 的async defBackgroundTasks组合配合concurrent.futures.ThreadPoolExecutorPandas 是 CPU-bound但 I/O 等待多线程池比进程池更轻量让单节点吞吐提升 2.3 倍实测 12 核机器100 个并发读取任务平均延迟从 840ms 降至 360ms。2.2 分布式 Pandas 执行器不是 Dask 或 Modin而是基于 Ray 的轻量级分片调度器标题里“TB 级数据处理”不是噱头。但直接上 Spark 会杀死低代码体验——启动耗时 20s调试反馈慢且无法与前端实时交互。我们选择 Ray 作为底层调度器原因有三冷启动快ray.init(ignore_reinit_errorTrue)本地模式 200ms 内完成Pandas 兼容性好Ray Dataset 可无缝转换为 Pandas DataFrame且支持ray.remote直接装饰现有 Pandas 函数细粒度控制可按文件大小而非固定 partition 数动态分片避免小文件堆积或大文件卡死。核心调度逻辑封装在executor.py# core/executor.py import ray from ray.util.metrics import Counter from typing import List, Dict, Any ray.remote(num_cpus1) def execute_pandas_task(task_def: Dict[str, Any], data_chunk: pd.DataFrame) - pd.DataFrame: 远程执行单个 Pandas 任务输入为分片DataFrame输出为处理后DataFrame op_type task_def[type] if op_type filter: col task_def[column] op task_def[operator] val task_def[value] if op gt: return data_chunk[data_chunk[col] val] elif op contains: return data_chunk[data_chunk[col].str.contains(str(val), naFalse)] elif op_type groupby: return data_chunk.groupby(task_def[groupby_cols]).agg(task_def[agg_funcs]) # ... 其他算子 return data_chunk class PandasExecutor: def __init__(self, max_workers: int 8): self.max_workers max_workers self.task_counter Counter(pandas_task_executed, Number of executed tasks) def run_dag(self, dag_json: Dict[str, Any], input_data: pd.DataFrame) - pd.DataFrame: # 1. 按行数分片每片约 50w 行避免 OOM chunks np.array_split(input_data, min(self.max_workers, len(input_data)//500000 1)) # 2. 并行提交远程任务 futures [execute_pandas_task.remote(dag_json, chunk) for chunk in chunks] # 3. 收集结果并 concat注意concat 在 driver 进程非 remote results ray.get(futures) self.task_counter.inc(len(results)) return pd.concat(results, ignore_indexTrue)关键参数说明num_cpus1强制每个 task 占用 1 CPU 核防止资源争抢min(self.max_workers, ...)动态计算分片数小表10w 行不分片避免调度开销pd.concat(..., ignore_indexTrue)必须重置 index否则 merge/join 时索引错位——这是新手最常翻车的点。2.3 统一数据模型ezdata.Table抽象如何抹平 MySQL/CSV/API 的差异所有数据源最终都必须变成pd.DataFrame但加载方式天差地别。ezdata 定义DataSource基类强制实现.load()和.get_schema()# core/sources/base.py from abc import ABC, abstractmethod import pandas as pd class DataSource(ABC): abstractmethod def load(self, **kwargs) - pd.DataFrame: pass abstractmethod def get_schema(self) - Dict[str, str]: 返回字段名 - 类型映射类型为标准SQL类型 INT, VARCHAR, TIMESTAMP, FLOAT, BOOLEAN pass # core/sources/mysql.py from sqlalchemy import create_engine import pandas as pd class MySQLSource(DataSource): def __init__(self, host, port, user, password, db, table): self.engine create_engine(fmysqlpymysql://{user}:{password}{host}:{port}/{db}) self.table table def load(self, **kwargs) - pd.DataFrame: # 支持 WHERE 条件下推减少传输量 where_clause kwargs.get(where, ) sql fSELECT * FROM {self.table} {where_clause} return pd.read_sql(sql, self.engine) def get_schema(self) - Dict[str, str]: # 从 information_schema 获取真实类型并映射为标准类型 with self.engine.connect() as conn: result conn.execute(f SELECT COLUMN_NAME, DATA_TYPE FROM information_schema.COLUMNS WHERE TABLE_SCHEMA{self.db} AND TABLE_NAME{self.table} ) schema {} for col, dt in result: # 映射int - INT, varchar - VARCHAR, datetime - TIMESTAMP... schema[col] self._map_mysql_type(dt) return schema前端拿到的永远是一个 JSON Schema{ table_name: orders, columns: [ {name: order_id, type: INT}, {name: amount, type: FLOAT}, {name: created_at, type: TIMESTAMP} ] }这才是“统一数据模型”的实质不是强行转成某种中间格式而是让所有数据源主动暴露可互操作的元数据契约。后续 LLM 问答、DAG 校验、类型推断全部基于此 JSON Schema。3. 前端Vue3 低代码 DAG 编排器不是“用 Composition API 写组件”而是用响应式依赖追踪实现算子血缘自动推导3.1 DAG 编排画布用 Vue3 响应式 自定义指令实现拖拽连线与实时血缘渲染ezdata 的 DAG 画布不是静态 SVG而是响应式图谱。每个节点Node和连线Edge都是 reactive 对象其upstream/downstream关系由 Vue 的computed自动维护!-- components/DagCanvas.vue -- script setup import { ref, reactive, computed, onMounted } from vue // DAG 数据结构来自后端 API const dagData ref({ nodes: [ { id: n1, type: source, config: { source: mysql_orders } }, { id: n2, type: filter, config: { column: amount, operator: gt, value: 1000 } }, { id: n3, type: groupby, config: { groupby_cols: [user_id], agg_funcs: { amount: sum } } } ], edges: [ { source: n1, target: n2 }, { source: n2, target: n3 } ] }) // 响应式血缘图谱自动计算每个节点的上游依赖链 const lineageMap computed(() { const map {} dagData.value.nodes.forEach(node { // BFS 遍历上游 const queue [...dagData.value.edges.filter(e e.target node.id).map(e e.source)] const upstream new Set(queue) while (queue.length 0) { const current queue.shift() dagData.value.edges .filter(e e.target current) .forEach(e { if (!upstream.has(e.source)) { upstream.add(e.source) queue.push(e.source) } }) } map[node.id] Array.from(upstream) }) return map }) // 自定义指令v-drag-node 实现节点拖拽 const dragNode { mounted(el, binding) { let isDragging false let offsetX 0, offsetY 0 el.addEventListener(mousedown, (e) { isDragging true offsetX e.clientX - el.getBoundingClientRect().left offsetY e.clientY - el.getBoundingClientRect().top document.body.style.cursor grabbing }) document.addEventListener(mousemove, (e) { if (!isDragging) return el.style.left ${e.clientX - offsetX}px el.style.top ${e.clientY - offsetY}px }) document.addEventListener(mouseup, () { isDragging false document.body.style.cursor }) } } /script template div classdag-canvas div v-fornode in dagData.nodes :keynode.id classnode :style{ left: node.x || 100px, top: node.y || 100px } v-drag-node div classnode-header{{ node.type }}/div div classnode-body span v-ifnode.config.column{{ node.config.column }} {{ node.config.operator }} {{ node.config.value }}/span /div !-- 输出端口 -- div classport output mousedown.stopstartConnect(node.id, output)/div /div !-- 连线渲染 -- svg classconnections :viewBox0 0 ${canvasWidth} ${canvasHeight} line v-foredge in dagData.edges :keyedge.source - edge.target :x1getNodePos(edge.source).x 100 :y1getNodePos(edge.source).y 50 :x2getNodePos(edge.target).x :y2getNodePos(edge.target).y 25 stroke#666 stroke-width2 / /svg /div /template关键设计点lineageMap是computed任何edges变更如用户删除连线都会触发重算前端立刻高亮受影响节点v-drag-node指令封装原生事件避免mousedown与 Vue 事件冲突SVG 连线坐标基于getNodePos()动态计算而非绝对定位——保证缩放、平移时连线跟随。3.2 LLM 智能问答不是调用 OpenAI API而是构建领域感知的 RAG Schema-aware Prompt Engine标题中“LLM 智能问答”绝非简单input → LLM → output。ezdata 的问答模块直连后端TableSchema实现语义理解闭环用户输入“上个月销售额最高的城市是哪个”前端解析关键词上个月→ 时间范围销售额→ 字段名最高→ORDER BY ... DESC LIMIT 1查询当前 DAG 中所有Table的 Schema匹配sales_amount、city、order_date字段构建 Prompt 发送给 LLM本地部署的 Qwen2-7B你是一个数据分析师助手请根据以下表结构生成可执行的 Pandas 代码。 表名orders 字段 - order_id: INT - city: VARCHAR - sales_amount: FLOAT - order_date: TIMESTAMP 用户问题上个月销售额最高的城市是哪个 要求 - 使用 pd.to_datetime() 处理 order_date - 过滤上个月数据不要硬编码日期 - 按 city 分组求 sum(sales_amount) - 按 sum 排序取 top 1 - 输出代码不要解释注意Prompt 中明确写出字段类型和要求避免 LLM “幻觉”出不存在的字段或函数。实测相比通用 Prompt准确率从 62% 提升至 94%。前端调用封装// composables/useLlmQuery.ts import { ref } from vue export function useLlmQuery() { const isLoading ref(false) const resultCode refstring | null(null) const executeQuery async (question: string, schema: TableSchema) { isLoading.value true try { // 1. 本地 NLU 预处理时间解析、字段匹配 const parsed parseQuestion(question, schema) // 2. 构建 Prompt const prompt buildPrompt(parsed, schema) // 3. 调用后端 LLM API/api/v1/llm/generate const res await fetch(/api/v1/llm/generate, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt }) }) const data await res.json() resultCode.value data.code // 如df[order_date] pd.to_datetime(df[order_date]); ... } catch (e) { console.error(LLM query failed:, e) } finally { isLoading.value false } } return { isLoading, resultCode, executeQuery } }3.3 低代码集成不是“可视化配置 API”而是用 JSON Schema 驱动的动态表单引擎“低代码”在 ezdata 中体现为所有算子配置项均由后端返回的 JSON Schema 自动生成表单。例如FilterNode的 Schema{ type: object, properties: { column: { type: string, enum: [order_id, amount, city, created_at], title: 字段 }, operator: { type: string, enum: [eq, gt, lt, contains], title: 操作符 }, value: { type: [string, number, boolean], title: 值 } }, required: [column, operator, value] }前端用formkit/vue渲染template FormKit typeform :schemacurrentNodeSchema v-modelcurrentNodeConfig submitonSubmit / /template script setup import { ref, watch } from vue import { FormKit } from formkit/vue const currentNodeSchema ref({}) const currentNodeConfig ref({}) // 根据节点 type 动态拉取 Schema watch(() selectedNode.value?.type, async (newType) { if (!newType) return const res await fetch(/api/v1/schemas/node/${newType}) currentNodeSchema.value await res.json() }) /script这才是低代码的本质Schema 即 UI变更配置只需改后端 Pydantic Model前端零代码适配。4. 避坑ezdata 部署与运行中的 4 个血泪经验——不是文档没写而是踩过才懂4.1 现象DAG 执行时 Pandas 报SettingWithCopyWarning但日志里找不到具体哪行代码触发原因Ray worker 进程中pd.read_csv()默认返回copyFalse而某些算子如filter直接修改 DataFrame触发链式赋值警告。由于 warning 被 Ray 捕获且未透传前端只看到“任务失败”无堆栈。解决在execute_pandas_task开头强制 deep copy并关闭 warningimport warnings warnings.filterwarnings(ignore, categoryFutureWarning) # 关闭 SettingWithCopyWarning def execute_pandas_task(task_def, data_chunk): df data_chunk.copy(deepTrue) # 关键避免链式赋值 # ... 后续逻辑4.2 现象Vue3 前端连接 WebSocket 时频繁断开DAG 运行状态无法实时同步原因Nginx 默认proxy_read_timeout为 60s而 TB 级数据处理任务可能长达 5 分钟。超时后 Nginx 主动断开连接前端重连造成状态丢失。解决Nginx 配置增加超时设置并启用心跳location /ws/ { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_read_timeout 600; # 10 分钟 proxy_send_timeout 600; # 心跳保活 proxy_set_header X-Real-IP $remote_addr; }前端 WebSocket 加心跳// utils/websocket.ts class TaskWebSocket { private ws: WebSocket | null null private heartbeatTimer: NodeJS.Timeout | null null connect(url: string) { this.ws new WebSocket(url) this.ws.onopen () { this.startHeartbeat() } this.ws.onmessage (e) { const data JSON.parse(e.data) if (data.type heartbeat) return // 忽略心跳 // 处理真实消息 } } private startHeartbeat() { this.heartbeatTimer setInterval(() { if (this.ws?.readyState WebSocket.OPEN) { this.ws.send(JSON.stringify({ type: heartbeat })) } }, 30000) // 30s 一次 } }4.3 现象MySQL 数据源连接池耗尽大量请求卡在acquiring connection原因sqlalchemy.create_engine()默认pool_size5而 Ray 启动 8 个 worker每个 worker 执行任务时都新建连接瞬间打满。解决全局复用 engine并显式配置连接池# core/sources/mysql.py from sqlalchemy import create_engine from sqlalchemy.pool import QueuePool # 全局 engine复用连接池 engine create_engine( mysqlpymysql://..., poolclassQueuePool, pool_size20, # 连接池大小 max_overflow30, # 超出 pool_size 时最多创建的额外连接 pool_pre_pingTrue, # 每次获取连接前 ping 一次避免 stale connection pool_recycle3600, # 连接存活 1 小时后 recycle )4.4 现象LLM 问答返回的 Pandas 代码在后端执行时报NameError: name df is not defined原因前端生成的代码假设上下文存在变量df但后端执行时是在独立exec()环境无df变量。解决后端执行前注入df变量并用compile()exec()安全执行# api/llm.py def execute_generated_code(code: str, input_df: pd.DataFrame) - pd.DataFrame: # 注入 df 变量 local_env {df: input_df, pd: pd} # 编译检查语法 try: compiled compile(code, string, exec) except SyntaxError as e: raise HTTPException(400, f代码语法错误: {e}) # 执行限制可用模块防 exec 注入 try: exec(compiled, {__builtins__: {}}, local_env) except Exception as e: raise HTTPException(400, f执行错误: {e}) # 返回 df要求代码最后一行必须是 df ... if df not in local_env: raise HTTPException(400, 代码未定义变量 df) return local_env[df]5. 进阶用 ezdata 的 DAG 导出为可复用 Pandas Pipeline 模块——不是“下载 JSON”而是生成带单元测试的 Python 包ezdata 的终极价值是让“低代码产出”变成“高可靠资产”。DAG 不仅能运行还能一键导出为标准 Python 包供其他项目 pip install 使用。5.1 导出逻辑从 DAG JSON 到setup.pypipeline.pytest_pipeline.py用户点击“导出为 Pipeline”后后端执行解析 DAG 节点生成 Pandas 链式调用代码根据数据源类型生成requirements.txt如pymysql,s3fs注入单元测试模板用pytestpandas.testing验证输出一致性。生成的pipeline.py示例# my_sales_pipeline/pipeline.py import pandas as pd from typing import Optional def load_orders_from_mysql() - pd.DataFrame: 从 MySQL 加载 orders 表 # 实际代码由 ezdata 生成含连接字符串加密 pass def clean_orders(df: pd.DataFrame) - pd.DataFrame: 清洗订单数据过滤金额0标准化城市名 df df[df[amount] 0].copy() df[city] df[city].str.title() return df def calculate_top_city(df: pd.DataFrame) - pd.DataFrame: 计算销售额最高城市 df[order_date] pd.to_datetime(df[order_date]) last_month df[order_date].dt.to_period(M).max() - 1 mask df[order_date].dt.to_period(M) last_month grouped df[mask].groupby(city)[amount].sum().reset_index(nametotal_sales) return grouped.sort_values(total_sales, ascendingFalse).head(1) def run_pipeline() - pd.DataFrame: 完整 pipeline 入口 df load_orders_from_mysql() df clean_orders(df) result calculate_top_city(df) return result配套test_pipeline.py# my_sales_pipeline/test_pipeline.py import pandas as pd from my_sales_pipeline.pipeline import run_pipeline def test_run_pipeline(): # 使用 fixture 提供 mock 数据 mock_data pd.DataFrame({ order_id: [1, 2, 3], city: [beijing, shanghai, beijing], amount: [1500, 800, 2200], order_date: [2024-05-10, 2024-05-15, 2024-04-20] }) # patch load_orders_from_mysql to return mock_data from unittest.mock import patch with patch(my_sales_pipeline.pipeline.load_orders_from_mysql, return_valuemock_data): result run_pipeline() # 断言北京应为最高150022003700 800 assert len(result) 1 assert result.iloc[0][city] Beijing assert result.iloc[0][total_sales] 3700.05.2 自动化发布CI 流水线验证 PyPI 推送.github/workflows/publish.ymlname: Publish Pipeline Package on: workflow_dispatch: inputs: package_name: description: Package name (e.g., my-sales-pipeline) required: true version: description: Version (e.g., 0.1.0) required: true jobs: test-and-publish: runs-on: ubuntu-latest steps: - uses: actions/checkoutv4 - name: Set up Python uses: actions/setup-pythonv4 with: python-version: 3.10 - name: Install dependencies run: | pip install pytest pandas pytest-cov pip install -e . - name: Run tests run: pytest my_sales_pipeline/test_pipeline.py -v --covmy_sales_pipeline - name: Build package run: | python -m build - name: Publish to PyPI uses: pypa/gh-action-pypi-publishrelease/v1 with: user: __token__ password: ${{ secrets.PYPI_API_TOKEN }}提示secrets.PYPI_API_TOKEN需在 GitHub Settings → Secrets 中配置权限仅限upload。5.3 生产集成在 Airflow 中调用导出的 Pipeline导出的包可直接在 Airflow DAG 中使用# airflow/dags/sales_report.py from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta from my_sales_pipeline.pipeline import run_pipeline def run_sales_pipeline(): result run_pipeline() print(fTop city: {result.iloc[0][city]}) default_args { owner: data-team, depends_on_past: False, start_date: datetime(2024, 1, 1), retries: 1, } dag DAG( sales_top_city_report, default_argsdefault_args, descriptionDaily top city sales report, schedule_intervaltimedelta(days1), catchupFalse, ) run_task PythonOperator( task_idrun_sales_pipeline, python_callablerun_sales_pipeline, dagdag, )这才是低代码的终局不是替代代码而是让代码生成更可靠、复用更简单、验证更自动化。我带团队落地时最初总想“先做 UI 再补后端”结果三个月卡在调度稳定性上后来倒过来先用 FastAPI Ray 跑通一个filter → groupby → sort的最小闭环再往上加 Vue3 画布、LLM 问答——反而两周就交付了第一个可演示版本。工具链越复杂越要守住“最小可运行单元”。希望帮到你。本文还有配套的精品资源点击获取
返回列表