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

如何定义Flink Schema以读取Pulsar中的Protocol Buffer数据

处理Pulsar-Flink中的Protocol Buffer数据

我之前在做Flink对接Pulsar处理Protobuf数据的时候也遇到过这个问题,官方确实没提供现成的Protobuf Schema实现,不过有两种很实用的解决办法,给你详细说下,附Scala代码示例:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 17:57:54