Spark任务正常但SPLINE UI无数据血缘,请求技术指导
问题描述
代码可正常完成数据读写,但Spline无法采集到元数据,在Spline UI(http://localhost:9090/app/events/list)中看不到任何数据。
环境版本:
- Spark 3.3.1
- Scala 2.12.18
- Python 3.9.6
- Spline Agent 1.1.0
最初通过spark-submit命令带参数提交时出错,改为在脚本内配置SparkConf后运行无报错,但Spline仍无数据产出。
原提交命令:
spark-submit --packages za.co.absa.spline.agent.spark:spark-3.3-spline-agent-bundle_2.12:1.1.0 --conf "spark.sql.queryExecutionListeners=za.co.absa.spline.harvester.listener.SplineQueryExecutionListener" --conf "spark.spline.producer.url=http://localhost:8080/producer" pyspark_example.py
当前脚本内配置:
conf = SparkConf().set("spark.sql.warehouse.dir", "./spark-warehouse").set("spark.jars.packages", "za.co.absa.spline.agent.spark:spark-3.3-spline-agent-bundle_2.12:1.1.0").set("spark.sql.queryExecutionListeners", "za.co.absa.spline.harvester.listener.SplineQueryExecutionListener").set("spark.spline.producer.url", "http://localhost:8080/producer")
完整PySpark代码:
from pyspark import SparkContext from pyspark.sql import SparkSession from pyspark.conf import SparkConf sc = SparkContext() conf = SparkConf() .set("spark.sql.warehouse.dir", "./spark-warehouse") .set("spark.jars.packages", "za.co.absa.spline.agent.spark:spark-3.3-spline-agent-bundle_2.12:1.1.0") .set("spark.sql.queryExecutionListeners", "za.co.absa.spline.harvester.listener.SplineQueryExecutionListener") .set("spark.spline.producer.url", "http://localhost:8080/producer") spark = SparkSession.builder.master("local[*]").appName("employee").config(conf = conf).getOrCreate() df = spark.read.csv("employee.csv") df.write.mode("overwrite").csv("sample")
解决方案
1. 修正Spark上下文初始化顺序
你代码里先创建了SparkContext,之后才配置SparkConf并创建SparkSession,这会导致Spline的配置无法加载到已初始化的Spark上下文里。SparkConf必须在SparkContext(或SparkSession)创建前设置,否则配置不会生效。
修改代码,先构建配置和SparkSession,无需单独创建SparkContext:
from pyspark.sql import SparkSession from pyspark.conf import SparkConf conf = SparkConf() .set("spark.sql.warehouse.dir", "./spark-warehouse") .set("spark.jars.packages", "za.co.absa.spline.agent.spark:spark-3.3-spline-agent-bundle_2.12:1.1.0") .set("spark.sql.queryExecutionListeners", "za.co.absa.spline.harvester.listener.SplineQueryExecutionListener") .set("spark.spline.producer.url", "http://localhost:8080/producer") # 直接通过SparkSession构建上下文,无需单独创建SparkContext spark = SparkSession.builder.master("local[*]").appName("employee").config(conf=conf).getOrCreate() df = spark.read.csv("employee.csv") df.write.mode("overwrite").csv("sample") # 显式停止会话,确保元数据发送完成 spark.stop()
2. 确认Spline服务状态
- 检查Spline Producer服务(http://localhost:8080)是否正常运行,可访问
http://localhost:8080/actuator/health查看健康状态 - 查看Spark应用日志,确认是否有连接Producer失败的报错信息
3. 启用Spline调试日志
在SparkConf中添加日志配置,查看Spline相关的详细日志,定位问题:
conf.set("log4j.logger.za.co.absa.spline", "DEBUG")
4. 确保操作被Spline追踪
Spline默认追踪SQL和结构化DataFrame操作,可尝试将DataFrame注册为临时表,用SQL语句执行读写,提升追踪成功率:
df = spark.read.csv("employee.csv").toDF("col1", "col2") # 建议指定列名,增强元数据可读性 df.createOrReplaceTempView("employee") spark.sql("INSERT OVERWRITE DIRECTORY 'sample' SELECT * FROM employee")
验证步骤
- 重启Spline Producer和UI服务
- 运行修改后的代码
- 查看Spark日志中是否有Spline相关的成功日志(如"Successfully persisted lineage")
- 刷新Spline UI查看元数据
内容的提问来源于stack exchange,提问作者user22032314
相关产品推荐
相关产品推荐

