【发布时间】:2017-11-01 03:50:13
【问题描述】:
我正在从 Kafka 获取数据,然后使用默认解码器反序列化 Array[Byte],之后我的 RDD 元素看起来像 (null,[B@406fa9b2)、(null,[B@21a9fe0) 但我想要我的原始数据有一个架构,那么我该怎么做实现这个?
我以 Avro 格式序列化消息。
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming spark-avro
我正在从 Kafka 获取数据,然后使用默认解码器反序列化 Array[Byte],之后我的 RDD 元素看起来像 (null,[B@406fa9b2)、(null,[B@21a9fe0) 但我想要我的原始数据有一个架构,那么我该怎么做实现这个?
我以 Avro 格式序列化消息。
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming spark-avro
您必须使用适当的反序列化器解码字节,例如字符串或您的自定义对象。
如果您不进行解码,您会得到 [B@406fa9b2,这只是 Java 中字节数组的文本表示。
Kafka 对消息的内容一无所知,因此它将字节数组从生产者传递给消费者。
在 Spark Streaming 中,您必须对键和值使用序列化程序(引用 KafkaWordCount example):
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer")
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer")
使用上述序列化程序,您将获得DStream[String],因此您可以使用RDD[String]。
但是,如果您想直接将字节数组反序列化为自定义类,则必须编写自定义 Serializer(这是 Kafka 特有的,与 Spark 无关)。
我建议使用带有固定架构的 JSON 或 Avro(使用Kafka, Spark and Avro - Part 3, Producing and consuming Avro messages 中描述的解决方案)。
在Structured Streaming 中,管道可能如下所示:
val fromKafka = spark.
readStream.
format("kafka").
option("subscribe", "topic1").
option("kafka.bootstrap.servers", "localhost:9092").
load.
select('value cast "string") // <-- conversion here
【讨论】:
map操作符。还没有from_avro(如果有的话),就像我们对带有from_json 的JSON 所做的那样。