
Spline架构全景图从Agent采集到ArangoDB存储的完整数据链路【免费下载链接】splineData Lineage Tracking And Visualization Solution项目地址: https://gitcode.com/gh_mirrors/spl/splineSpline 是一个开源的数据血缘追踪与可视化解决方案专为 Apache Spark 等数据计算框架设计。一句话概括它的工作Spark Agent 把作业的执行计划 执行事件上报给 Spline ServerServer 经过模型映射后通过 FOxx 事务把血缘关系写成 ArangoDB 里的图结构再由消费端 REST API 输出血缘概览、字段级血缘和影响分析结果。本文带你完整走一遍这条Spline 数据链路看清每个模块的职责与边界。一、整体架构速览一条主链路两类通道 整个系统可以抽象成一条单向数据流 一个查询出口Spark Agent血缘采集端 │ 执行计划 / 执行事件JSON ▼ Spline Server ├── REST 通道producer-rest-core └── Kafka 通道kafka-gateway │ ModelMapper 版本翻译 幂等校验 ▼ FOxx 事务arangodb-foxx-services ▼ ArangoDB图存储节点 边集合 ▲ Consumer REST APIconsumer-rest-core └── 血缘概览 / 详细血缘 / 影响分析两个关键设计决策双通道接入REST 简单直接Kafka 适合高吞吐、削峰和容错图数据库存储血缘天然是一张谁读谁、谁写谁的有向图ArangoDB 的图模型正好匹配查询血缘路径可以用图遍历高效完成。二、采集端血缘数据从哪里来 ️血缘数据不是 Spline Server 主动爬取的而是由部署在 Spark 集群侧的 Agent 拦截产生Agent 在 Spark 作业提交时捕获执行计划Operation、Attribute、Expression、DataSource运行结束后捕获执行事件进度、错误、耗时。Agent 支持两种上报协议对应服务端两个入口模块通道服务端模块适用场景RESTproducer-rest-core/中小规模、快速接入Kafkakafka-gateway/大规模集群、削峰、断点重试三、REST 通道多版本 API 的接入与兼容REST 入口的核心是 IngestionController.scala完整路径producer-rest-core/src/main/scala/za/co/absa/spline/producer/rest/controller/IngestionController.scala。它通过混入HandlerV10、HandlerV11、HandlerV12三个 Handler同时兼容三代 API 版本——老版本 Agent 上报的 JSON 也能被正确解析这是 Spline 平滑升级的关键。请求进来后还经过一层过滤器链在rest-gateway/的启动类中注册GzipFilter压缩传输减小大计划体的网络开销MessageLengthCapturingFilter记录请求体大小便于排查消息过长类问题。启动配置见rest-gateway/src/main/scala/za/co/absa/spline/gateway/rest/AppInitializer.scala它把producer上报、consumer查询、about诊断三组 REST 服务挂在同一个应用里一套部署同时承担写入与读取。四、Kafka 通道高吞吐异步接入的可靠性设计Kafka 通道的核心是 IngestionListener.scala路径kafka-gateway/src/main/scala/za/co/absa/spline/gateway/kafka/listener/IngestionListener.scala它订阅指定 topic把消息交给与 REST 通道完全相同的 Repository 落库——两条通道在存储层汇合。根据kafka-gateway/README.md这套消费者有几个值得学习的可靠性设计幂等同一条消息消费多次不会产生重复数据Spline 识别到已入库即忽略重试 退避失败消息按指数退避重试可通过spline.kafka.backOff.*系列参数调节死信队列DLQ开启spline.kafka.deadLetterQueueEnabledtrue后只有数据成功入库或已转入 DLQ 才提交 offset保证不丢数据水平扩展多实例并行消费或用spline.kafka.consumerConcurrency在单实例内扩容。五、模型映射ModelMapper 如何抹平版本差异不同 API 版本的 JSON 字段结构不同统一转换由producer-model-mapper/完成。其抽象接口非常简洁trait ModelMapper[P, E] { def fromDTO(plan: P): ExecutionPlan def fromDTO(event: E): ExecutionEvent }对应路径producer-model-mapper/src/main/scala/za/co/absa/spline/producer/modelmapper/ModelMapper.scala。三个版本各有一个实现v1_0/ModelMapperV10.scala面向最古老的 Spark 旧格式内部还通过spark/子包做属性依赖解析v1_1/ModelMapperV11.scala、v1_2/ModelMapperV12.scala面向新格式映射逻辑更直接。映射产物统一为内部规范模型producer-model/下的ExecutionPlan、ExecutionEvent下游服务因此完全不用关心版本差异——这是典型的防腐层设计。六、生产者服务幂等写入与 UUID 冲突检测拿到规范模型后由producer-services/模块负责落库。核心类是 ExecutionProducerRepositoryImpl.scala路径producer-services/src/main/scala/za/co/absa/spline/producer/service/repo/ExecutionProducerRepositoryImpl.scala每次写入都做了三件事查重先用 AQL 按_key查询是否已存在同名计划/事件若存在则比对discriminator冲突检测ID 相同但 discriminator 不一致时抛出UUIDCollisionDetectedException防止两个不同作业因 UUID 巧合互相污染血缘重试整个写入过程包裹在AsyncCallRetryer中遇到可重试的数据库异常自动重试。写入计划前还会先把计划引用的数据源用 AQLUPSERT方式去重入库按uri幂等保证 DataSource 节点全局唯一。七、存储层FOxx 事务如何把血缘写进 ArangoDB 图 ️这是整条链路最有意思的一层。Spline 没有用应用侧的多步写入而是把一次计划入库要写十几张表封装成数据库内的一个原子事务——由 ArangoDB FOxx 服务承载。应用侧persistence/模块的FoxxPostTxBuilder把持久化模型 POST 到 FOxx 端点/spline/execution-plans或/spline/execution-events数据库侧arangodb-foxx-services/中的 plans-router.ts路径arangodb-foxx-services/src/main/routes/plans-router.ts接收请求调用execution-plan-store.ts执行存储。真正写入时storeExecutionPlan在单个事务内依次插入节点集合与边集合失败则整体回滚节点集合边集合血缘关系ExecutionPlan、Operation、Schema、Attribute、Expression、DataSourceExecutes / Depends / Affects、ReadsFrom / WritesTo、Emits / Uses / Produces、ConsistsOf、ComputedBy / DerivesFrom、Takes、Follows这些边就是数据血缘的本体ReadsFrom表示操作读自数据源DerivesFrom表示字段派生自字段ComputedBy表示字段由表达式计算。查询端正是沿着这些边做图遍历才能回答某字段从哪来、某张表影响谁。八、消费端血缘查询的 REST 出口 写入完成后查询侧由两个模块协同consumer-rest-core/REST 控制器层每个业务查询一个 Controller——LineageOverviewController血缘概览图的全景LineageDetailedController详细血缘字段级路径ImpactOverviewController影响分析写操作波及哪些下游另有OperationDetailsController、LabelsController、DataSourcesController等。consumer-services/仓储层LineageRepository、ImpactRepository等把查询翻译成对 ArangoDB 图边集合的遍历。对外还暴露了分页能力Pageable、PageRequest模型支撑 UI 上按页浏览作业列表这类交互。九、运维支撑Admin CLI 与数据库迁移 数据库初始化/升级admin/模块提供 AdminCLI命令定义在admin/src/main/scala/za/co/absa/spline/admin/commands.scala包括DBInit建库建集合、DBUpgrade执行迁移、DataRetention数据保留策略清理等另通过arango/foxx/FoxxManager.scala管理 FOxx 服务的部署与卸载迁移脚本编排persistence/src/main/scala/za/co/absa/spline/persistence/migration/MigrationScriptRepository.scala把所有版本对(vFrom → vTo)建成一张图用最短路算法自动拼出从任意旧版本到最新版本的迁移链——比如 0.1 直升到 3.0也能自动找到 0.1→1.0→…→3.0 的脚本序列版本检查persistence/中的DatabaseVersionManager在应用启动时校验数据库 schema 版本与应用版本是否匹配避免新旧不兼容的隐性故障。十、快速上手从克隆到跑通完整链路 本地体验整条数据链路只需三步# 1. 克隆源码 git clone https://gitcode.com/gh_mirrors/spl/spline # 2. 构建需 Java 11 Maven 3.6 cd spline mvn install # 3. 启动 ArangoDB 与 Spline 服务后Kafka 通道最小配置示例 # -Dspline.database.connectionUrlarangodb://localhost/spline # -Dspline.kafka.consumer.bootstrap.serverslocalhost:9092 # -Dspline.kafka.topicspline-topic配置项细节可参考kafka-gateway/README.mdDocker 部署则在构建时追加-Ddocker参数即可生成镜像。总结一张图看懂 Spline 的架构精髓回顾整条链路Spline 的四个亮点值得借鉴双通道接入 统一存储层REST 与 Kafka 入口殊途同归可靠性能力幂等、重试、DLQ只在一处实现ModelMapper 防腐层三代 API 版本互不干扰升级无痛数据库内事务一次血缘写入涉及十几张集合FOxx 单事务保证要么全成、要么全无图模型存储血缘即图边集合命名即语义查询端直接图遍历。理解了Agent → 双通道 → 映射 → 事务落图 → 图遍历查询这条主线再阅读producer-rest-core/、kafka-gateway/、arangodb-foxx-services/、consumer-services/四个核心模块的源码就会顺畅得多。【免费下载链接】splineData Lineage Tracking And Visualization Solution项目地址: https://gitcode.com/gh_mirrors/spl/spline创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考