SparkSession集成Pub/Sub报错:数据源pubsub未找到,如何解决?
解决Spark连接Pub/Sub时数据源未找到的问题
核心问题分析
你遇到的DATA_SOURCE_NOT_FOUND错误主要源于两个关键问题:
- 依赖包不匹配:你使用的
spark-streaming-pubsub是针对旧版Spark Streaming(DStream API)的依赖,而代码采用的是Spark SQL/DataFrame API,需要对应SQL版本的连接器包。 - 读写模式错误:Pub/Sub连接器主要支持流处理模式,你的代码用
spark.read(批量读取)尝试加载数据,会导致数据源识别失败。
具体修复步骤
1. 替换正确的依赖包
修改SparkSession配置中的spark.jars.packages,使用与你的Spark版本匹配的SQL连接器包:
- 若Spark版本为2.4.x(Scala 2.11):
.config("spark.jars.packages", "org.apache.bahir:spark-sql-pubsub_2.11:2.4.0") - 若Spark版本为3.x(Scala 2.12):
.config("spark.jars.packages", "org.apache.bahir:spark-sql-pubsub_2.12:3.2.0")
2. 切换为流处理模式
Pub/Sub连接器专为流处理设计,将代码中的批量读写改为流处理模式:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession \ .builder \ .appName("PubSubSpark") \ .config("spark.jars.packages", "org.apache.bahir:spark-sql-pubsub_2.11:2.4.0") \ .getOrCreate() project_number = "********" topic = "*****" sub = "*****" # 从Pub/Sub读取流数据 stream_df = spark.readStream \ .format("pubsub") \ .option("subscriptionId", f"projects/{project_number}/subscriptions/{sub}") \ .option("topicId", f"projects/{project_number}/topics/{topic}") \ .load() # 测试:将流数据输出到控制台 query = stream_df.writeStream \ .outputMode("append") \ .format("console") \ .start() # 保持流处理运行 query.awaitTermination() # 若需写入到另一个Pub/Sub Topic,可使用以下代码 # write_query = stream_df.writeStream \ # .format("pubsub") \ # .option("topicId", f"projects/{project_number}/topics/{target_topic}") \ # .option("checkpointLocation", "/tmp/pubsub_checkpoint") # 必须指定检查点路径 # .start() # write_query.awaitTermination() spark.stop()
3. 额外排查点
- 若网络受限无法自动下载依赖包,可手动下载对应版本的jar包,放置到Spark安装目录的
jars文件夹下。 - 确认Spark版本、Scala版本与依赖包版本完全匹配,避免版本兼容问题。
- 检查GCP权限配置,确保Spark应用拥有访问Pub/Sub Topic和订阅的权限。
内容的提问来源于stack exchange,提问作者Chidananda Nayak
相关产品推荐
相关产品推荐

