如何在Spark中使用ScalaPB运行时泛型解包google.protobuf.Any类型
问题原因
你遇到的编译错误和需求实现难点核心有两个层面:
- 调用
unpack传入动态Class时,返回值类型是无边界的java.lang.Object,Spark无法为泛化的Any类型自动生成序列化所需的Encoder - 运行时解包需要确保对应类型的Proto描述符/实体类在Executor的类路径中可用
解决方案
分两种常用场景提供实现方式:
场景1:需要拿到具体的Proto实体类实例,后续用通用Message接口处理
显式指定返回值类型为com.google.protobuf.Message,并手动提供对应的Kryo序列化Encoder即可解决编译报错,示例代码:
import scalapb.spark.Implicits._ import org.apache.spark.sql.Encoders import com.google.protobuf.Message import java.util.Base64 // 显式定义Message类型的序列化编码器 implicit val messageEncoder: org.apache.spark.sql.Encoder[Message] = Encoders.kryo(classOf[Message]) // 把加载的类强转为Message的子类 val entityClass = Class.forName(qualifiedEntityClassNameFromRuntime).asSubclass(classOf[Message]) spark.read.format("text") .load("/tmp/entity").as[String] .map { s => Base64.getDecoder.decode(s) } .map { bytes => EntityEnvelope.parseFrom(bytes) } .map { envelope => toJavaProto(envelope.getEntity).unpack(entityClass) } .show()
场景2:不需要依赖具体实体类,仅需要读取实体字段内容
可以直接用DynamicMessage基于描述符解析,不需要提前知道实体类名,直接从Any的typeUrl提取类型信息即可:
import scalapb.spark.Implicits._ import com.google.protobuf.DynamicMessage import java.util.Base64 // 提前维护所有可能出现的实体类型的Descriptor映射,key为实体的全限定名 val descriptorMap: Map[String, com.google.protobuf.Descriptors.Descriptor] = Map( "your.package.EntityABC" -> EntityABC.javaDescriptor, "your.package.EntityXYZ" -> EntityXYZ.javaDescriptor ) spark.read.format("text") .load("/tmp/entity").as[String] .map { s => Base64.getDecoder.decode(s) } .map { bytes => val envelope = EntityEnvelope.parseFrom(bytes) val javaAny = toJavaProto(envelope.getEntity) // 从typeUrl中提取实体的全限定名 val fullTypeName = javaAny.getTypeUrl.substring(javaAny.getTypeUrl.lastIndexOf('/') + 1) val descriptor = descriptorMap(fullTypeName) // 直接基于描述符解析出动态消息,不需要依赖具体实体类 DynamicMessage.parseFrom(descriptor, javaAny.getValue) } .show()
如果需要得到结构化的DataFrame,可以在map中用com.google.protobuf.util.JsonFormat.printer()把DynamicMessage转成JSON字符串,后续直接用Spark的JSON数据源解析即可,不需要提前定义Schema。
内容的提问来源于stack exchange,提问作者Mike Dias
相关产品推荐
相关产品推荐

