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

PySpark向Google Pub/Sub Lite发消息遇NoSuchMethodError求助

问题分析与解决方案

这个错误是Guava版本冲突导致的:Spark默认自带的旧版Guava与Pub/Sub Lite依赖的高版本Guava存在方法签名不兼容,具体是com.google.common.base.Preconditions.checkArgument方法的重载版本在旧版中不存在,引发NoSuchMethodError。

解决步骤

1. 调整依赖参数,排除Spark自带的Guava并指定兼容版本

修改PYSPARK_SUBMIT_ARGS,添加Guava的依赖并排除Spark内置的旧版本:

os.environ['PYSPARK_SUBMIT_ARGS'] = '''--packages com.google.cloud:pubsublite-spark-sql-streaming:0.4.1,com.google.cloud:google-cloud-pubsublite:1.6.1,com.google.guava:guava:30.1.1-jre 
--exclude-packages com.google.guava:guava 
pyspark-shell'''

2. 验证依赖版本匹配

  • pubsublite-spark-sql-streaming:0.4.1适配的Guava版本为30.x系列,避免使用过高或过低版本
  • 必须排除Spark自带的guava包,强制使用指定的高版本,否则旧版会覆盖依赖的高版本

3. 额外注意事项

  • 若使用Spark 2.x版本,可能需要将Guava版本调整为28.2-jre,需根据实际环境兼容性测试
  • 确保checkpointLocation路径唯一且具备读写权限,避免重复运行时的状态冲突

修正后的完整代码片段(关键部分)

import os
import uuid
from pyspark.sql import SparkSession
from pyspark.sql.functions import array, create_map, col, lit, when
from pyspark.sql.types import BinaryType, StringType

# 配置依赖,解决Guava版本冲突
os.environ['PYSPARK_SUBMIT_ARGS'] = '''--packages com.google.cloud:pubsublite-spark-sql-streaming:0.4.1,com.google.cloud:google-cloud-pubsublite:1.6.1,com.google.guava:guava:30.1.1-jre 
--exclude-packages com.google.guava:guava 
pyspark-shell'''

# TODO(developer): 替换为你的项目信息
project_number = xxx
location = "us-central1"
topic_id = "kosmin"

spark = SparkSession.builder.appName("write-app").getOrCreate()

# 假设sdf是已定义的输入数据流
sdf = (
    sdf.withColumn("key", lit("example").cast(BinaryType()))
        .withColumn("data", col("value").cast(StringType()).cast(BinaryType()))
        .withColumnRenamed("timestamp", "event_timestamp")
        .withColumn(
            "attributes",
            create_map(
                lit("key1"),
                array(when(col("value") % 2 == 0, b"even").otherwise(b"odd")),
            ),
        )
        .drop("value")
)

query = (
    sdf.writeStream.format("pubsublite")
        .option(
            "pubsublite.topic",
            f"projects/{project_number}/locations/{location}/topics/{topic_id}",
        )
        .option("checkpointLocation", "/tmp/app" + uuid.uuid4().hex)
        .outputMode("append")
        .start()
)

query.awaitTermination(60)
query.stop()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 14:06:19