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

PySpark结构化流中基于protoc生成类读取Protobuf消息的问题

Spark结构化流读取Protobuf消息:binaryDescriptorSet方案可行性分析

我通过protoc生成了对应消息类型的Python类,在Spark(Kafka -> 结构化流)中读取Protobuf消息时,发现了from_protobuf函数,但该函数需要我未持有的descFilePath,或是从现有类中难以获取的binaryDescriptorSet。我尝试了以下代码,想确认这种通过binaryDescriptorSet的方案是否可行:

decoded_df = encoded_df.select(
    from_protobuf(
        encoded_df.value,
        messageName=Message().DESCRIPTOR.name,
        binaryDescriptorSet=bytearray.fromhex(
            str(Message().DESCRIPTOR.file.serialized_pb)
        ),
    ))

如果该方案不可行,我计划创建UDF将Protobuf消息转换为Spark Struct。另外我了解到pbspark,但它不符合我的项目要求。


方案可行性分析及修正

你当前的binaryDescriptorSet用法存在错误:serialized_pb本身就是bytes类型的二进制描述符集合,不需要转成十六进制字符串再解码,直接传递即可。

修正后的binaryDescriptorSet用法

decoded_df = encoded_df.select(
    from_protobuf(
        encoded_df.value,
        messageName=Message().DESCRIPTOR.name,
        binaryDescriptorSet=Message().DESCRIPTOR.file.serialized_pb
    ))

需要注意的前提:如果你的Protobuf消息依赖其他外部.proto文件(即关联消息定义分散在多个文件),则需要将所有依赖文件的serialized_pb合并为一个字节数组后再传递,Spark需要完整的描述符集合才能正确解析嵌套或关联的消息结构。

若上述方案失效:UDF替代方案

如果描述符集合的处理仍有问题,UDF是更稳妥的选择,完全基于你已有的Python Protobuf类实现:

  1. 定义与Protobuf消息结构匹配的Spark StructType
  2. 编写UDF完成二进制消息的解析与格式转换

示例代码:

from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

# 定义与Protobuf消息对应的Spark Schema
message_schema = StructType([
    StructField("id", IntegerType(), nullable=False),
    StructField("content", StringType(), nullable=True)
])

def parse_protobuf(msg_bytes):
    try:
        msg = Message()
        msg.ParseFromString(msg_bytes)
        return (msg.id, msg.content)
    except Exception:
        return (None, None)

# 注册UDF
parse_protobuf_udf = udf(parse_protobuf, message_schema)

# 解析消息
decoded_df = encoded_df.select(parse_protobuf_udf(encoded_df.value).alias("parsed_msg"))

UDF方案的优势是灵活适配现有类,无需处理描述符集合的复杂逻辑;缺点是性能略低于原生from_protobuf,大流量场景下需注意资源消耗。

内容的提问来源于stack exchange,提问作者Francisco Fonseca

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 09:25:05