SQL 血缘支持
当我们执行一条 SQL 语句时,计算引擎会返回一个结果集,结果集中的每一列可能来源于不同的表,而这些不同的表又依赖于其他表,这些表可能经过了函数、聚合等计算。源表与结果集之间的关系,就是 SQL 血缘(SQL Lineage)。
简介
当前的血缘解析功能是通过扩展 Spark 的 QueryExecutionListener 以插件形式实现的。
- 在 SQL 执行完成后会触发
SparkListenerSQLExecutionEnd事件,并由QueryExecuctionListener捕获,随后对执行成功的 SQL 进行血缘解析处理。 - 解析得到的血缘信息会以 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
AlterViewAsCommandAppendDataCreateDataSourceTableAsSelectCommandCreateHiveTableAsSelectCommandCreateTableAsSelectCreateViewCommandInsertIntoDataSourceCommandInsertIntoDataSourceDirCommandInsertIntoHadoopFsRelationCommandInsertIntoHiveDirCommandInsertIntoHiveTableMergeIntoTableOptimizedCreateHiveTableAsSelectCommandOverwriteByExpressionOverwritePartitionsDynamicReplaceTableAsSelectSaveIntoDataSourceCommand
构建
使用 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.x | x | - |
| 3.0.x | x | - |
| 2.4.x | x | - |
目前支持基于 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.address | http://localhost:21000 | Atlas 服务器的 REST 端点 URL | 1.8.0 |
| atlas.client.type | rest | 客户端类型(目前仅支持 rest) | 1.8.0 |
| atlas.client.username | none | 客户端用户名 | 1.8.0 |
| atlas.client.password | none | 客户端密码 | 1.8.0 |
| atlas.cluster.name | primary | 实体 qualifiedName 中使用的集群名称 | 1.8.0 |
| atlas.hook.spark.column.lineage.enabled | true | 是否将字段级血缘导入 Atlas | 1.8.0 |
评论
登录后参与评论
KnowForge