You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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")
验证步骤
  1. 重启Spline Producer和UI服务
  2. 运行修改后的代码
  3. 查看Spark日志中是否有Spline相关的成功日志(如"Successfully persisted lineage")
  4. 刷新Spline UI查看元数据

内容的提问来源于stack exchange,提问作者user22032314

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 18:33:15