欢迎回来
登录你的知识库账户
忘记密码?
还没有账户?立即注册
创建账户
注册你的专属知识库
已有账户?去登录
找回密码
输入注册邮箱获取验证码
返回登录
请输入图片中的验证码以继续注册
加载中...
取消
新建收藏
手动添加你喜欢的内容
取消
编辑头像与昵称
上传新头像或修改你的显示昵称
支持 JPG/PNG,最大 2MB
取消

问题反馈

notebasewww.notebase.cn
控制台
内容库
动态
管理
账户
U
用户
--
在线
v0.8.7 · 知识库
笔记
KnowledgeBase
网络无边,知识有迹。
0笔记
0工具
30推荐

分类导航

按主题直达

编辑精选

站内用户贡献 · 真实笔记

最新收录

每日更新
继续浏览全部内容 →
>
笔记
0
加载中...
工具
0
此页用于记录用户反馈问题后的每一次改进
笔记用法

“写笔记”支持四种格式——Word 文档、Excel 表格、Markdown、纯文本,起稿或二次编辑时都能随时切换,同一篇笔记想用哪种形态来记,都由你说了算。

md、txt、csv、json 这类纯文本则原样载入,不做多余加工。拿一张现成的表倒进来、改几笔、再导出去,等于白用一台免费的格式转换器。

要带走就在右上角点“下载”,可导出 PDF、Word、Markdown、Excel、TXT 等格式;列表卡片“⋯”菜单里,也有同样的下载入口。

工具用法

在“工具”页点“+ 上传工具”即可发布:填好名称与链接,再用 Markdown 把使用方法写清楚——能解决什么问题、怎么装、怎么用,比堆介绍实在。

要分发安装包就一并上传压缩包(ZIP、RAR、7Z、TAR.GZ,最大 35MB),别人在详情页一键下载;只放链接不带附件也可以。

工具按大家的收藏热度排序,好用的自然会被顶上来。发布后可在详情页或卡片菜单里编辑、下架。

隐藏笔记

写笔记时勾上“隐藏”,这篇就只存在于你自己的账号里:不进列表、不进搜索、不上首页精选,也不会出现在任何公开的页面,链接发给别人同样打不开。

适合放密码、草稿、日记这类只给自己看的内容;想公开,去“发布”打开它,把“隐藏”的勾去掉再保存,之后编辑会默认保持原状态,不会悄悄变回公开。

数据安全

你的内容会同时保存在多个副本上,系统定期做备份与完整性校验,再配合异地容灾机制:就算某台机器出问题,数据也不会丢,可以长期放心存放;特别重要的资料,仍建议你另外再留一份备份。

技术

全站跑在容器化、模块化的现代架构上,更新、部署、回滚都很快,扩展性和稳定性都按长期运营的标准来设计(Built for reliability, designed to scale)。

理念

这个网站最早只是一个人的笔记仓库,后来慢慢长成现在的知识中枢。设计上很克制——没有广告、没有追踪、没有推荐算法,只是干干净净地存放一些东西;既然做好了,就公开出来,万一有人用得上呢。

原则

不做大而全,不做平台梦,保持简单、保持克制、保持好奇。所有内容都由用户贡献、由用户维护:不会突然冒出付费墙,不会在角落塞广告位,也不会把你的数据卖给第三方。

更多

产品会持续迭代,站内日志页记录着每一次改动,改了什么都有迹可循;想了解这个站是怎么一步步走到今天的,翻翻日志就能看到来龙去脉。

举报

如果在这里看到涉嫌违规的内容,点对应卡片右侧的“举报”按钮就能提交,我们会尽快核实处理;也谢谢你花一点时间,一起把这里维护干净。

趋势
// 点击导航加载发现
归档
// 归档为空
最近浏览
// 暂无浏览记录
发布
// 加载中...
用户发布
// 加载中...
用户管理
// 加载中...
访问统计
// 加载中...
内容审核
// 加载中...
个人信息
// 加载中...
返回首页

Spark App 血缘解析方案

2022/9/16技术教程

背景

数据量越来越大之后,数据血缘这东西就变得越来越重要了。不管是追查问题数据的来源,还是发现那些没人用的“孤岛”表想下线,都得靠血缘来梳理表跟表、表跟任务、任务跟任务之间的关系。

我们之前已经用 ANTLR 做了 SQL 任务的血缘解析,但 Spark App 这块一直靠人工配置,太原始了。所以这次想把这个补上,把 Spark App 的血缘也自动解析出来。

现状是这样的:

  • 线上跑的 Spark App 有 Spark 2.3 和 Spark 3.1 两个版本
  • 支持的语言包括 Python 2/3、Java、Scala
  • 运行平台有 YARN 和 K8s
  • 血缘收集机制需要能覆盖上面所有情况

思路

Spark App 的血缘解析,大致有三条路可以走:

1. 基于代码解析

直接解析 Spark App 的源代码来提取血缘关系。类似的产品有 SPROV。
但问题是 Spark App 写法太灵活了,要同时支持 Java、Python、Scala 三种语言,复杂度直接爆炸,不太现实。

2. 基于动态监听

运行时通过修改代码或者插件的方式收集血缘。典型的有:

  • Titian、Pebble — 需要改代码
  • Spline、Apache Atlas — 插件方式,侵入性小

3. 基于日志解析

分析 Spark App 的 event log,从中提取血缘信息。

我们最开始觉得日志解析这条路比较省事,就去翻了 Spark 2 和 Spark 3 的历史 event log。结果发现:

Spark 2 的 event log 里没有完整的 Hive 表元信息,而 Spark 3 在 FileSourceScanExec、HiveTableScan 这些算子中打出了 Hive 表信息。

所以基于 event log 的方式没法完美支持 Spark 2,这条路就堵死了。

最终我们决定走 动态监听 这条路,调研了 Spline,做了可用性分析。下面详细说。


Spline 调研

Spline(Spark Lineage)是一个 Apache 2.0 协议开源的 Spark 血缘收集系统,免费。整体分三部分:

  • Spline Agent — 负责血缘解析,核心部分
  • Spline Server — 负责接收和存储血缘数据
  • Spline UI — 展示界面

因为我们内部有自己的血缘展示系统,所以主要关注的是 Spline Agent 这块,Server 和 UI 可以用我们自己的替代。

架构图

Spark App → Spline Agent → LineageDispatcher → Spline Server → 存储/展示

初始化方式

Spline 支持两种初始化方式,本质都是注册一个 QueryExecutionListener 来监听 SparkListenerSQLExecutionEnd 消息。

1. Codeless 初始化(推荐)

不需要改代码,通过配置直接嵌入。启动命令示例:

spark-submit \
  --jars /path/to/lineage/spark-3.1-spline-agent-bundle_2.12-1.0.0-SNAPSHOT.jar \
  --files /path/to/lineage/spline.properties \
  --num-executors 2 \
  --executor-memory 1G \
  --driver-memory 1G \
  --name test_lineage \
  --deploy-mode cluster \
  --conf spark.spline.mode=BEST_EFFORT \
  --conf spark.spline.lineageDispatcher.http.producer.url=http://172.18.221.156:8080/producer \
  --conf "spark.sql.queryExecutionListeners=za.co.absa.spline.harvester.listener.SplineQueryExecutionListener" \
  test.py

关键点:

  • --jars 指定 spline agent 的 jar 包(也可以放到 Spark 部署的 jars 目录下)
  • --files 指定 spline 配置文件,或者直接用 --conf 传配置(注意要加 spark. 前缀)
  • spark.sql.queryExecutionListeners 注册监听器

2. Programmatic 初始化

需要在代码里显式调用。支持 Scala、Java、Python。

Scala 示例:

val sparkSession: SparkSession = ???

import za.co.absa.spline.harvester.SparkLineageInitializer._
sparkSession.enableLineageTracking()

Java 示例:

import za.co.absa.spline.harvester.SparkLineageInitializer;

SparkLineageInitializer.enableLineageTracking(session);

Python 示例:

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("spline_app") \
    .config("spark.jars", "dbfs:/path_where_the_jar_is_uploaded") \
    .getOrCreate()

sc = spark.sparkContext
sc.setSystemProperty("spline.mode", "REQUIRED")
jvm = sc._jvm
jvm.za.co.absa.spline.harvester.SparkLineageInitializer.enableLineageTracking(spark._jsparkSession)

血缘解析原理

核心逻辑在 SplineAgent.handle() 方法里,最终调用 LineageHarvester.harvest() 来获取血缘。

通过 SparkListenerSQLExecutionEnd 消息拿到 QueryExecution,然后基于其中的:

  • analyzed logical plan
  • executedPlan

进行解析。

解析流程

  1. 解析写操作:tryExtractWriteCommand(logicalPlan)
    通过插件机制,从 PluginRegistry 中获取 WriteNodeProcessing 类型的插件,解析出写操作对应的 Hive 表信息。
    比如 DataSourceV2Plugin.writeNodeProcessor() 专门处理 V2WriteCommand、CreateTableAsSelect、ReplaceTableAsSelect 这些命令。

    插件可以自己扩展,继承 za.co.absa.spline.harvester.plugin.Plugin 就行,Spline Agent 启动时自动加载 classpath 里的所有插件。

  2. 解析读操作:基于写操作中的 query 字段,递归解析 logical plan,找出所有读取的表。

  3. 输出:最终得到两个 JSON:

    • plan — 血缘关系
    • event — 辅助信息

输出示例

{
  "plan": {
    "id": "acd5157c-ddc5-5ef0-b1bc-06bb8dcda841",
    "name": "team evaluation ranks",
    "operations": {
      "write": {
        "outputSource": "hdfs:///user/hive/warehouse/dm_ai.db/dws_kdt_comment_ranks_info",
        "append": false,
        "id": "op-0",
        "name": "CreateDataSourceTableAsSelectCommand",
        "childIds": ["op-1"],
        "params": {
          "table": {
            "identifier": {
              "table": "dws_kdt_comment_ranks_info",
              "database": "dm_ai"
            },
            "storage": "Storage()"
          }
        },
        "extra": {
          "destinationType": "orc"
        }
      },
      "reads": [
        {
          "inputSources": [
            "hdfs://yz-cluster-qa/user/hive/warehouse/dm_ai.db/dws_kdt_comment_rank_base"
          ],
          "id": "op-6",
          "name": "LogicalRelation",
          "output": [
            "attr-0", "attr-1", "attr-2", "attr-3", "attr-4",
            "attr-5", "attr-6", "attr-7", "attr-8", "attr-9",
            "attr-10", "attr-11", "attr-12"
          ],
          "params": {
            "table": {
              "identifier": {
                "table": "dws_kdt_comment_rank_base",
                "database": "dm_ai"
              },
              "storage": "Storage(Location: hdfs://yz-cluster-qa/user/hive/warehouse/dm_ai.db/dws_kdt_comment_rank_base, Serde Library: org.apache.hadoop.hive.ql.io.orc.OrcSerde, InputFormat: org.apache.hadoop.hive.ql.io.orc.OrcInputFormat, OutputFormat: org.apache.hadoop.hive.ql.io.orc.OrcOutputFormat, Storage Properties: [serialization.format=1])"
            }
          },
          "extra": {
            "sourceType": "hive"
          }
        }
      ],
      "other": [
        {
          "id": "op-5",
          "name": "SubqueryAlias",
          "childIds": ["op-6"],
          "output": [
            "attr-0", "attr-1", "attr-2", "attr-3", "attr-4",
            "attr-5", "attr-6", "attr-7", "attr-8", "attr-9",
            "attr-10", "attr-11", "attr-12"
          ],
          "params": {
            "identifier": "spark_catalog.dm_ai.dws_kdt_comment_rank_base"
          }
        },
        {
          "id": "op-4",
          "name": "Filter",
          "childIds": ["op-5"],
          "output": [
            "attr-0", "attr-1", "attr-2", "attr-3", "attr-4",
            "attr-5", "attr-6", "attr-7", "attr-8", "attr-9",
            "attr-10", "attr-11", "attr-12"
          ],
          "params": {
            "condition": {
              "__exprId": "expr-0"
            }
          }
        }
      ]
    }
  }
}

可以看到,写操作里能拿到目标表的 database、table、outputSource,读操作里能拿到输入表的完整信息,包括 Hive 元数据。这些信息足够我们构建完整的血缘关系了。

总结

  • Event log 方案:Spark 2 不支持,放弃
  • 代码解析方案:语言太多太复杂,放弃
  • Spline 方案:基于动态监听,支持 Spark 2.3 和 3.1,支持 Java/Scala/Python,YARN/K8s 都能跑,插件可扩展,目前看是最靠谱的

接下来就是实际集成到我们的任务调度系统里,把 Spline Agent 默认加到每个 Spark App 的启动参数中,然后对接上我们自己的血缘存储和展示系统。

编写使用方法
Markdown 格式 · Ctrl+Enter 确定
新建笔记
预览
数据表格
点击单元格编辑 · Tab 移动
A1fx
Sheet1
BIH1H2≡🔗</>
隐私提醒

取消
编辑工具
取消