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

SparkSession集成Pub/Sub报错:数据源pubsub未找到,如何解决?

解决Spark连接Pub/Sub时数据源未找到的问题

核心问题分析

你遇到的DATA_SOURCE_NOT_FOUND错误主要源于两个关键问题:

  1. 依赖包不匹配:你使用的spark-streaming-pubsub是针对旧版Spark Streaming(DStream API)的依赖,而代码采用的是Spark SQL/DataFrame API,需要对应SQL版本的连接器包。
  2. 读写模式错误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 04:16:27