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类实现:
- 定义与Protobuf消息结构匹配的Spark
StructType - 编写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
相关产品推荐
相关产品推荐

