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

使用Delta表时PySpark StreamingQueryListener的QueryTerminatedEvent未触发问题

Delta流中StreamingQueryListener的QueryTerminatedEvent触发失败及报错解决

问题分析

从错误日志来看,核心问题是Py4J通信异常,包含连接被拒绝、空命令错误两类。对比可用的CSV示例和Delta示例代码,问题根源集中在以下几点:

  • 重复创建SparkSession:Delta代码中两次调用configure_spark_with_delta_pip(builder).getOrCreate(),导致Spark上下文实例混乱,Py4J回调通道绑定到失效实例,引发通信故障。
  • 未正确管理Listener生命周期:CSV示例在流程结束前移除了Listener,但Delta代码未执行此操作,当Query终止时,Listener尝试通过已失效的Py4J通道传递事件,触发连接拒绝。
  • Query终止流程不完整:调用q.stop()后未等待Query完全终止,程序直接结束,导致QueryTerminatedEvent无法正常传递到Listener。

解决方案

  1. 单例SparkSession:整个程序只创建一次SparkSession,避免上下文实例冲突。
  2. 规范Listener生命周期:在Query完全终止后,手动移除Listener,切断与失效上下文的绑定。
  3. 等待Query终止完成:调用q.awaitTermination()确保Query有足够时间完成终止流程,让Listener能正常接收终止事件。
  4. 修正路径错误:原代码中output_table =" spark/output"存在多余空格,会导致路径识别异常,需修正为正确格式。

修正后的Delta示例代码

from pyspark.sql.streaming import StreamingQueryListener
from pyspark.sql.streaming.listener import (
    QueryStartedEvent,
    QueryProgressEvent,
    QueryTerminatedEvent,
)
import pyspark
import time
from delta import configure_spark_with_delta_pip

# 仅初始化一次SparkSession
builder = (
    pyspark.sql.SparkSession.builder.appName("MyApp")
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
    .config(
        "spark.sql.catalog.spark_catalog",
        "org.apache.spark.sql.delta.catalog.DeltaCatalog",
    )
)
spark = configure_spark_with_delta_pip(builder).getOrCreate()

input_table = "spark/input"
output_table = "spark/output"  # 修正原代码的路径空格问题
mydata1 = "tmp/data_dump/part-00000-17f59480-16e2-418a-960e-285fdd04ee45.c000.snappy.parquet"
mydata2 = "tmp/data_dump/part-00000-248f4fb4-7c85-4b13-baaf-0689b57e54c2.c000.snappy.parquet"

class CustomListener(StreamingQueryListener):
    def onQueryStarted(self, event):
        print("STARTED")

    def onQueryProgress(self, event):
        print("PROGRESS")

    def onQueryTerminated(self, event):
        print("ENDED")

# 添加自定义Listener
mylistener = CustomListener()
spark.streams.addListener(mylistener)

# 创建源和目标Delta表
batch_df = spark.read.format("parquet").load(mydata1)
batch_df.write.format("delta").save(input_table)
batch_df.write.format("delta").save(output_table)

# 启动流查询
df = spark.readStream.format("delta").load(input_table)
q = (df.writeStream
    .option("checkpointLocation", "checkpoint")
    .outputMode("append")
    .format("delta")
    .trigger(processingTime="1 seconds")
    .start(output_table)
)

time.sleep(5)
print("extra data")
batch_df2 = spark.read.format("parquet").load(mydata2)
batch_df2.write.format("delta").mode("append").save(input_table)
time.sleep(5)

# 停止查询并等待终止完成
q.stop()
q.awaitTermination(timeout=10)  # 最多等待10秒确保终止流程完成

# 移除Listener并关闭SparkSession
spark.streams.removeListener(mylistener)
spark.stop()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 07:09:54