Spark 扩展

SQL 血缘支持

qianmoQqianmoQ· 更新于 2026-09-29· 阅读 11 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

当我们执行一条 SQL 语句时,计算引擎会返回一个结果集,结果集中的每一列可能来源于不同的表,而这些不同的表又依赖于其他表,这些表可能经过了函数、聚合等计算。源表与结果集之间的关系,就是 SQL 血缘(SQL Lineage)。

简介

当前的血缘解析功能是通过扩展 Spark 的 QueryExecutionListener 以插件形式实现的。

  1. 在 SQL 执行完成后会触发 SparkListenerSQLExecutionEnd 事件,并由 QueryExecuctionListener 捕获,随后对执行成功的 SQL 进行血缘解析处理。
  2. 解析得到的血缘信息会以 JSON 格式写入日志文件。

示例

执行以下 SQL 语句时:

## table
create table test_table0(a string, b string)

## query
select a as col0, b as col1 from test_table0

该 SQL 的血缘:

{
   "inputTables": ["spark_catalog.default.test_table0"],
   "outputTables": [],
   "columnLineage": [{
      "column": "col0",
      "originalColumns": ["spark_catalog.default.test_table0.a"]
   }, {
      "column": "col1",
      "originalColumns": ["spark_catalog.default.test_table0.b"]
   }]
}

血缘特定标识

  • __count__。表示该列是一个 count(*) 聚合表达式,无法提取具体列。类似 default.test_table0.__count__ 这样的列的血缘。
  • __local__。表示该表的血缘是一个 LocalRelation,而非真实表,例如 __local__.a

SQL 类型支持

目前支持 Spark 的 Command 和 Query 类型的列级血缘:

Query

  • Select

Command

  • AlterViewAsCommand
  • AppendData
  • CreateDataSourceTableAsSelectCommand
  • CreateHiveTableAsSelectCommand
  • CreateTableAsSelect
  • CreateViewCommand
  • InsertIntoDataSourceCommand
  • InsertIntoDataSourceDirCommand
  • InsertIntoHadoopFsRelationCommand
  • InsertIntoHiveDirCommand
  • InsertIntoHiveTable
  • MergeIntoTable
  • OptimizedCreateHiveTableAsSelectCommand
  • OverwriteByExpression
  • OverwritePartitionsDynamic
  • ReplaceTableAsSelect
  • SaveIntoDataSourceCommand

构建

使用 Apache Maven 构建

Kyuubi Spark Lineage Listener Extension 使用 Apache Maven 构建。要构建它,请 cd 进入 kyuubi 项目的根目录并运行:

build/mvn clean package -pl :kyuubi-spark-lineage_2.12 -am -DskipTests

稍后,如果一切顺利,你最终会得到两部分插件产物:

  • 主插件 jar,位于 ./extensions/spark/kyuubi-spark-lineage/target/kyuubi-spark-lineage_${scala.binary.version}-${project.version}.jar

针对不同 Apache Spark 版本构建

Maven 选项 spark.version 用于指定编译所用的 Spark 版本,并生成相应的传递依赖。默认情况下,始终使用 Kyuubi 项目主 pom 文件中定义的最新 spark.version 进行构建。有时它可能与其他 Spark 发行版不兼容,此时你可能需要针对自己所使用的 Spark 版本自行构建该插件。

例如,

build/mvn clean package -pl :kyuubi-spark-lineage_2.12 -am -DskipTests -Dspark.version=3.5.1

可用的 spark.version 见下表。

Spark 版本是否支持备注
master√-
3.5.x√-
3.4.x√-
3.3.x√-
3.2.x√-
3.1.xx-
3.0.xx-
2.4.xx-

目前支持基于 Scala 2.12 发布的 Spark。

使用 ScalaTest Maven 插件进行测试

如果在上面的命令中省略 -DskipTests 选项,你还将运行所有的单元测试。

build/mvn clean package -pl :kyuubi-spark-lineage_2.12

如果出现任何 bug,且你希望自己调试该插件,可以配置 -DdebugForkedProcess=true,并可选地配置 -DdebuggerPort=5005(可选)。

build/mvn clean package -pl :kyuubi-spark-lineage_2.12 -DdebugForkedProcess=true

测试将在启动时暂停,并等待远程调试器连接到配置的端口。

如果您能将该问题或修复方案分享给 Kyuubi 社区,我们将不胜感激。

安装

确保 kyuubi-spark-lineage_*.jar 及其传递依赖位于 Spark 运行时的类路径中,例如:

  • 复制到 $SPARK_HOME/jars 目录,或
  • 通过 spark.jars 配置项进行指定

配置

Spark Listener 扩展配置

将 org.apache.kyuubi.plugin.lineage.SparkOperationLineageQueryExecutionListener 添加到 Spark 配置 spark.sql.queryExecutionListeners 中。

spark.sql.queryExecutionListeners=org.apache.kyuubi.plugin.lineage.SparkOperationLineageQueryExecutionListener

可选配置

是否跳过永久视图解析

如果启用,血缘解析将在永久视图处停止,并将其视为物理表。我们需要添加一个配置项。

spark.kyuubi.plugin.lineage.skip.parsing.permanent.view.enabled=true

获取血缘事件

血缘派发器用于派发血缘事件,通过 spark.kyuubi.plugin.lineage.dispatchers 进行配置。

  • SPARK_EVENT(默认):将血缘事件发送到 Spark 事件总线
  • KYUUBI_EVENT:将血缘事件发送到 Kyuubi 事件总线
  • ATLAS:将血缘发送到 Apache Atlas

从 SparkListener 获取血缘事件

使用 SPARK_EVENT 派发器时,血缘事件会被发送到 SparkListenerBus。要处理血缘事件,需要添加一个新的 SparkListener。添加 SparkListener 的示例如下:

spark.sparkContext.addSparkListener(new SparkListener {
      override def onOtherEvent(event: SparkListenerEvent): Unit = {
        event match {
          case lineageEvent: OperationLineageEvent =>
            // Your processing logic
          case _ =>
        }
      }
    })

从 Kyuubi EventHandler 获取血缘事件

使用 KYUUBI_EVENT 分发器时,血缘事件将被发送到 Kyuubi 的 EventBus。有关如何处理 Kyuubi 事件,请参阅 Kyuubi Event Handler。

将血缘实体导入 Apache Atlas

使用 ATLAS 分发器可以将血缘实体导入 Apache Atlas。

额外的准备工作:

  • 所需的最小传递依赖位于 ./extensions/spark/kyuubi-spark-lineage/target/scala-${scala.binary.version}/jars
  • 使用 spark.files 指定 Atlas 的 atlas-application.properties 配置文件

Atlas 客户端配置(在 atlas-application.properties 中配置,或通过 spark.atlas. 前缀传入):

名称默认值说明引入版本
atlas.rest.addresshttp://localhost:21000Atlas 服务器的 REST 端点 URL1.8.0
atlas.client.typerest客户端类型(目前仅支持 rest)1.8.0
atlas.client.usernamenone客户端用户名1.8.0
atlas.client.passwordnone客户端密码1.8.0
atlas.cluster.nameprimary实体 qualifiedName 中使用的集群名称1.8.0
atlas.hook.spark.column.lineage.enabledtrue是否将字段级血缘导入 Atlas1.8.0

评论

登录后参与评论

正在加载评论…