PySpark Streaming处理Kafka Protobuf报bytearray无WhichOneof属性错误
错误根因
报错本质是类型不匹配:Spark通过Kafka连接器读入的value字段默认是原始二进制字节数组(bytearray类型),代码中没有做任何Protobuf反序列化操作,就直接把字节数组传入UDF,试图调用Protobuf对象专属的WhichOneof方法,自然会抛出属性不存在的错误。
快速排查
- 在UDF第一行加
print(type(msg)),启动任务后查看executor日志,会看到打印的类型是<class 'bytearray'>,而非预期的schema_pb2.MarketDataEvent实例,直接验证入参未被反序列化。 - 检查Kafka读入逻辑:当前代码仅指定了format为kafka,没有配置任何Protobuf反序列化规则,Spark不会自动将二进制字节转换为自定义Protobuf类。
解决方案
根据场景二选一即可:
方案1:UDF内手动反序列化(改动最小,适合本地调试/小流量场景)
直接修改UDF逻辑,先把传入的字节数组反序列化为Protobuf对象,再执行后续oneof字段判断,修改后的代码如下:
def parse_options_monitor_msg(msg_bytes: bytearray): # 先完成二进制到Protobuf对象的反序列化 msg = schema_pb2.MarketDataEvent() msg.ParseFromString(bytes(msg_bytes)) # 原有解析逻辑保留 eventStr = msg.WhichOneof("event") if eventStr == "trade": trade_msg = msg.trade symbol = trade_msg.symbol ticker = symbol.split('_')[0] return str([symbol, ticker]) # 补充oneof其他分支的返回逻辑,避免返回null elif eventStr == "instrumentStatusUpdate": return str(["instrument_status", "update"]) else: return str(["unknown", "event"]) # 去掉多余的lambda包装,直接传入函数即可 parse_options_monitor = udf(parse_options_monitor_msg, StringType()) # 后续流处理逻辑无需改动 df = spark.readStream \ .format("kafka") \ .options(**kafka_conf) \ .load() data = df.selectExpr("offset", "value") \ .withColumn("event", parse_options_monitor(col("value"))) df2 = data.select(col("offset"),col("event")) df2.writeStream \ .format("console") \ .outputMode("append") \ .option("truncate", False) \ .start() \ .awaitTermination()
方案2:使用Spark内置Protobuf函数反序列化(性能更优,适合生产环境)
Python UDF反序列化需要在Python和JVM之间做数据拷贝,性能差,Spark 3.4+版本内置了Protobuf解析支持,序列化/反序列化逻辑在JVM层执行,性能是Python UDF的数倍,实现步骤:
- 用protoc命令把proto文件编译成描述符集文件:
protoc --descriptor_set_out=market_data.desc --include_imports 你的proto文件路径.proto - 把生成的
market_data.desc放到所有Spark节点可访问的路径(本地路径、HDFS、对象存储路径均可) - 修改流处理代码,替换自定义UDF为内置解析函数:
from pyspark.sql.functions import from_protobuf, col, locate, substr df = spark.readStream \ .format("kafka") \ .options(**kafka_conf) \ .load() # 直接用内置函数完成二进制反序列化 data = df.selectExpr("offset", "value") \ .withColumn("parsed_event", from_protobuf( col("value"), "MarketDataEvent", descFilePath="market_data.desc的存放路径" ) ) # 直接从解析后的结构体中取字段,不需要写UDF df2 = data.select( col("offset"), col("parsed_event.trade.symbol").alias("symbol"), substr(col("parsed_event.trade.symbol"), 1, locate("_", col("parsed_event.trade.symbol"))-1).alias("ticker") ) df2.writeStream \ .format("console") \ .outputMode("append") \ .option("truncate", False) \ .start() \ .awaitTermination()
注意事项
- 用UDF反序列化的场景,必须保证集群所有Executor节点安装的protobuf Python版本,和生成
schema_pb2.py文件使用的protoc版本完全一致,否则会出现反序列化失败、字段错乱的问题。 - 用内置
from_protobuf的场景,如果Spark版本低于3.4,需要手动引入对应版本的spark-protobuf依赖包,才能使用相关函数。 - 无论用哪种方案,都要覆盖Protobuf中oneof定义的所有分支,避免未匹配分支时返回空值,增加后续排查成本。
内容的提问来源于stack exchange,提问作者aiman
相关产品推荐
相关产品推荐

