使用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。
解决方案
- 单例SparkSession:整个程序只创建一次SparkSession,避免上下文实例冲突。
- 规范Listener生命周期:在Query完全终止后,手动移除Listener,切断与失效上下文的绑定。
- 等待Query终止完成:调用
q.awaitTermination()确保Query有足够时间完成终止流程,让Listener能正常接收终止事件。 - 修正路径错误:原代码中
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
相关产品推荐
相关产品推荐

