如何定义Flink Schema以读取Pulsar中的Protocol Buffer数据
处理Pulsar-Flink中的Protocol Buffer数据
我之前在做Flink对接Pulsar处理Protobuf数据的时候也遇到过这个问题,官方确实没提供现成的Protobuf Schema实现,不过有两种很实用的解决办法,给你详细说下,附Scala代码示例:
方案一:自定义Flink DeserializationSchema
Flink的DeserializationSchema是用来定义如何将字节流转化为业务对象的接口,我们可以自己实现这个接口来处理Protobuf的反序列化逻辑。
假设你已经通过protoc生成了Protobuf实体类MyMessage,先实现自定义Schema:
import org.apache.flink.api.common.serialization.DeserializationSchema import org.apache.flink.api.common.typeinfo.TypeInformation import org.apache.flink.util.Collector import com.google.protobuf.InvalidProtocolBufferException class ProtobufSchema extends DeserializationSchema[MyMessage] { // 核心反序列化逻辑:将Pulsar传来的字节数组解析为Protobuf对象 override def deserialize(message: Array[Byte]): MyMessage = { try { MyMessage.parseFrom(message) } catch { case e: InvalidProtocolBufferException => // 根据业务需求处理异常:比如抛出异常终止任务,或者返回null跳过该条数据 throw new RuntimeException("Failed to parse Protobuf message", e) } } // 标记是否到达流末尾,Protobuf流一般不需要这个,直接返回false override def isEndOfStream(nextElement: MyMessage): Boolean = false // 返回生成对象的类型信息,Flink需要这个来做类型推断 override def getProducedType: TypeInformation[MyMessage] = { TypeInformation.of(classOf[MyMessage]) } }
然后在创建FlinkPulsarSource时使用这个自定义Schema:
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment import org.apache.flink.streaming.connectors.pulsar.FlinkPulsarSource import java.util.Properties object PulsarProtobufJob { def main(args: Array[String]): Unit = { val env = StreamExecutionEnvironment.getExecutionEnvironment val serviceUrl = "pulsar://your-pulsar-broker:6650" // 替换成你的Pulsar服务地址 val adminUrl = "http://your-pulsar-admin:8080" // 替换成你的Pulsar Admin地址 val props = new Properties() props.setProperty("topic", "test-source-topic") props.setProperty("subscription-name", "flink-protobuf-sub") // 必须设置订阅名称 // 使用自定义Schema创建Pulsar Source val source = new FlinkPulsarSource[MyMessage]( serviceUrl, adminUrl, new ProtobufSchema(), props ) // 添加Source并处理数据 val messageStream = env.addSource(source) messageStream.print() // 示例:打印Protobuf对象 env.execute("Flink-Pulsar Protobuf Processing Job") } }
方案二:复用Pulsar原生ProtobufSchema
Pulsar客户端本身已经提供了ProtobufSchema实现,我们可以通过Flink-Pulsar连接器的PulsarDeserializationSchema来包装它,这样不需要自己写解析逻辑,更简洁:
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment import org.apache.flink.streaming.connectors.pulsar.{FlinkPulsarSource, PulsarDeserializationSchema} import org.apache.pulsar.client.api.schema.ProtobufSchema import java.util.Properties object PulsarNativeProtobufJob { def main(args: Array[String]): Unit = { val env = StreamExecutionEnvironment.getExecutionEnvironment val serviceUrl = "pulsar://your-pulsar-broker:6650" val adminUrl = "http://your-pulsar-admin:8080" val props = new Properties() props.setProperty("topic", "test-source-topic") props.setProperty("subscription-name", "flink-protobuf-sub") // 复用Pulsar原生的ProtobufSchema val pulsarProtobufSchema = ProtobufSchema.of(classOf[MyMessage]) // 包装成Flink能识别的DeserializationSchema val flinkSchema = PulsarDeserializationSchema.pulsarSchema(pulsarProtobufSchema) val source = new FlinkPulsarSource[MyMessage]( serviceUrl, adminUrl, flinkSchema, props ) env.addSource(source).print() env.execute("Flink-Pulsar Native Protobuf Job") } }
注意事项
- Protobuf类生成:确保你已经通过protoc(或ScalaPB插件)正确生成了
MyMessage类,并且项目中引入了Protobuf的依赖(比如com.google.protobuf:protobuf-java) - 依赖版本匹配:Flink-Pulsar连接器的版本要和你的Flink、Pulsar版本保持兼容,避免出现依赖冲突
- 异常处理:在反序列化失败时,根据业务场景选择是跳过消息还是终止任务,避免因为脏数据导致整个任务崩溃
- 订阅配置:务必设置
subscription-name,否则Flink会使用默认订阅名,可能导致重复消费或消费异常
内容的提问来源于stack exchange,提问作者xksa
相关产品推荐
相关产品推荐

