公司动态

FlowLoom:基于 Apache Spark 的可视化数据处理平台

📅 2026/8/17 22:20:04
FlowLoom:基于 Apache Spark 的可视化数据处理平台
文章目录人工备注说明一、项目简介二、架构设计2.1 flow-model — 数据模型层2.2 flow-client — 可视化设计器2.3 flow-script — Spark 执行引擎2.4 子进程通信协议三、扩展机制3.1 自定义 Map 函数3.2 自定义 Filter 函数四、技术栈总览五、设计亮点六、总结七、运行效果7.1 使用的测试数据7.2 配置文件数据源节点7.3 配置Spark Map节点7.4 配置Spark Filter节点7.5 配置Spark SQL节点7.6 配置DB导出节点7.7 执行完成后数据表的数据人工备注说明开发工具Claude Code 火山方舟glm-5.1/deepseek v4 pro这个项目完全使用Claude Code开发大概2天左右完成大概2小时完成70%代码剩下30%代码经过多次测试和提示词进行微调完成。这篇文章也是使用Claude Code编写注意这个应用只进行简单测试未覆盖所有场景可能存在异常的情况文末提供源码供学习研究一、项目简介FlowLoom 是一个可视化数据处理桌面工具旨在让数据分析师和开发人员无需编写复杂的 Spark 代码通过拖拽方式即可设计并执行数据处理流程。核心工作流用户通过 Swing GUI 设计数据处理流程图 → 保存为 XML 文件 → Apache Spark 引擎分布式执行。支持的数据处理能力类别节点类型说明数据源文件数据源支持 Excel / CSV / TXT 文件读取数据源数据库数据源支持 MySQL、Oracle 等关系型数据库转换Map 映射通过自定义 Java 类实现字段级转换转换Filter 过滤按条件筛选数据行转换SQL 查询使用 SQL 语句进行数据查询与聚合输出CSV 导出结果导出为 CSV 文件输出Excel 导出结果导出为 Excel 文件输出数据库导出结果写入 MySQL 等数据库表二、架构设计FlowLoom 采用 Maven 多模块架构分为三个子模块依赖关系清晰flow-model共享数据模型无 UI/Spark 依赖 ↑ ↑ flow-client flow-script (Swing GUI) (Spark 引擎)2.1 flow-model — 数据模型层定义流程图的领域模型是 client 和 script 之间唯一的共享契约。核心类包括FlowGraph流程图顶层容器包含节点列表和边列表内置拓扑排序算法FlowNode流程节点包含类型、坐标位置和属性键值对FlowEdge有向边连接源节点和目标节点FlowProperty节点属性的键值对模型NodeType枚举所有节点类型分为 SOURCE / TRANSFORM / SINK 三大类别所有模型类使用 Jackson XML 注解进行序列化生成的 XML 文件兼具可读性和可移植性flowname数据处理流程typeFILEversion1.0nodesnodeidn1typeFILE_SOURCEx70y110propertynamefilePathvaluedata/input.csv/propertynamefileTypevalueCSV//nodenodeidn2typeFILTERx270y110propertynamefilterClassvalueTestFilter//nodenodeidn3typeCSV_EXPORTx470y110propertynamefilePathvaluedata/output.csv//node/nodesedgesedgeide1sourcen1targetn2/edgeide2sourcen2targetn3//edges/flow2.2 flow-client — 可视化设计器基于 Java Swing 构建的桌面应用技术选型精巧FlatLaf现代化 Swing 主题提供原生感的界面体验JGraphX (mxGraph)流程图绘制引擎支持拖拽、连线、缩放等交互RSyntaxTextArea代码编辑器组件用于 SQL 和自定义代码编辑zt-exec子进程管理库用于启动 flow-script 执行引擎界面布局采用经典的分割面板设计MainFrame ├── FlowMenuBar菜单栏 ├── TaskListPanel左侧任务列表 └── JTabbedPane右侧多标签页编辑器 └── TaskTabPanel单个任务编辑器 ├── NodeToolbarPanel / ActionToolbarPanel工具栏 ├── FlowCanvasPanel流程画布 │ ├── FlowGraphAdapter模型桥接层 │ ├── FlowGraphComponentmxGraph 组件 │ └── FlowGraphStyles样式定义 ├── PropertiesPanel属性面板CardLayout 切换 │ ├── FileSourcePropertiesPanel │ ├── MapPropertiesPanel │ ├── FilterPropertiesPanel │ ├── SqlPropertiesPanel │ ├── CsvExportPropertiesPanel │ ├── ExcelExportPropertiesPanel │ └── DbExportPropertiesPanel └── ExecutionLogPanel运行日志面板其中FlowGraphAdapter是关键的桥接层负责将FlowGraph模型对象与 mxGraph 可视化组件进行双向同步——用户在画布上的拖拽、连线、属性编辑操作都会实时反映到模型中反之亦然。2.3 flow-script — Spark 执行引擎执行引擎是整个平台的计算核心采用经典的管道Pipeline模式MainSparkFlow → FlowParser解析 XML 为 FlowGraph → PipelinePlan拓扑排序生成执行计划 → PipelineExecutor遍历执行计划 → FlowOperator算子分发执行FlowParser负责验证并解析 XML 流程文件生成FlowGraph领域对象。PipelinePlan对FlowGraph执行拓扑排序计算出无环依赖的执行顺序。如果流程图中存在循环依赖会直接抛出异常阻止执行。PipelineContext作为运行时上下文持有 SparkSession 和各节点输出的 DataFrame节点间通过 context 传递数据。PipelineExecutor内部维护一个EnumMapNodeType, FlowOperator的算子注册表按执行顺序依次调用对应算子publicclassPipelineExecutor{privatefinalMapNodeType,FlowOperatoroperatorsnewEnumMap(NodeType.class);publicPipelineExecutor(){registerOperator(newFileSourceOperator());registerOperator(newDbSourceOperator());registerOperator(newMapOperator());registerOperator(newFilterOperator());registerOperator(newSqlOperator());registerOperator(newCsvExportOperator());registerOperator(newExcelExportOperator());registerOperator(newDbExportOperator());}publicvoidexecute(PipelinePlanplan,PipelineContextctx)throwsException{for(FlowNodenode:plan.getExecutionOrder()){FlowOperatoroperatoroperators.get(node.getType());operator.execute(ctx,node);}}}所有算子实现统一的FlowOperator接口——仅两个方法getSupportedType()声明支持的节点类型execute()执行具体逻辑。这种设计使得新增算子类型只需实现接口并注册即可完全符合开闭原则。2.4 子进程通信协议flow-client 通过ProcessBuilder启动 flow-script 作为子进程执行。两者通过 stdout 的按行文本协议进行通信协议消息含义[LOG] message普通日志消息[ERROR] message错误消息[PROGRESS] nodeId percent节点执行进度[COMPLETE]流程执行成功[FAILED] message流程执行失败flow-script 的 Log4j2 配置中专门设置了 Protocol appender使用%msg%n格式确保协议消息不带时间戳前缀地输出到 stdout供 flow-client 的ProcessOutputConsumer解析并展示在日志面板中。三、扩展机制3.1 自定义 Map 函数Map 节点通过mapClass属性指定自定义转换类。用户只需继承MapFormat基类并实现转换逻辑publicclassDateFormatMapextendsMapFormat{OverridepublicRowmap(Rowrow){// 对每一行数据进行自定义转换returnrow;}}运行时通过反射实例化com.penngo.flowloom.script.map.{className}实现了插件化的扩展机制。3.2 自定义 Filter 函数类似 MapFilter 节点通过filterClass属性指定过滤类继承FilterFormat或使用模板类TemplateFilter在 GUI 中直接编写过滤表达式。四、技术栈总览技术版本用途Java17主语言Apache Spark4.1.1分布式计算引擎Jackson XML2.20.0XML 序列化Apache POI5.4.1Excel 读写FlatLaf—Swing 现代化主题JGraphX—流程图可视化Log4j22.20.0日志框架Lombok1.18.42减少样板代码MySQL Connector8.3.0数据库连接Hutool5.8.44通用工具库LangChain4j1.12.2AI 能力集成预留五、设计亮点模型驱动架构flow-model 作为单一数据源GUI 和引擎共享同一套模型定义避免了数据不一致问题算子插件化通过FlowOperator接口 EnumMap注册表实现算子的热插拔新增节点类型无需修改执行器核心代码拓扑排序执行自动计算 DAG 的执行顺序并检测循环依赖保证流程的正确性进程隔离执行GUI 与 Spark 引擎运行在不同 JVM 进程中通过文本协议通信避免了 Spark 依赖污染 GUI 进程也便于资源隔离和崩溃恢复编码兼容处理针对 Spark CSV 不支持 GBK 的问题FileSourceOperator 会自动检测编码并进行转码属性泛化设计节点属性以ListFlowProperty键值对存储而非强类型字段保证 XML 格式的向前兼容性和可扩展性六、总结FlowLoom 是一个设计精巧的数据处理平台它巧妙地将低代码可视化理念与Apache Spark 分布式计算能力结合起来。通过清晰的模块划分、插件化的算子体系和简洁的子进程通信协议在保持架构灵活性的同时降低了使用门槛。无论是日常的数据 ETL 任务还是需要快速搭建数据处理管道的场景FlowLoom 都能提供一个直观高效的解决方案。七、运行效果7.1 使用的测试数据7.2 配置文件数据源节点7.3 配置Spark Map节点7.4 配置Spark Filter节点7.5 配置Spark SQL节点7.6 配置DB导出节点7.7 执行完成后数据表的数据FlowLoom源码