【发布时间】:2020-02-21 08:59:24
【问题描述】:
我有一个 DStream[String,String]。我使用 foreachRDD 获取每个 RDD 并将其发布到 Kafka 中。 我遇到的问题是我需要保证 String 被序列化,并且我的 RDD 的值由于未知原因而无法序列化。 Kafka 期望将 StringSerializer 作为值,但正如您在下图中看到的那样,我的 DStream 没有序列化字符串。如何在发布到 Kafka 之前将不可序列化的 String 转换为 serializabel?我可以更改 kafConf,但我更喜欢更改值而不是 Kafka 配置。
def kafkaConf(brokers : String) = {
val props = new HashMap[String, Object]()
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers)
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")
props
}
【问题讨论】:
-
你为什么不分享你制作人的代码?
-
Spark Streaming 已弃用。你为什么不使用结构化流?并且请不要将图片用于错误
-
很抱歉没有提供有关我的错误的更多详细信息。我从老师那里收到了关于问题所在的答复。问题在于火花上下文。我将 spark 上下文作为参数传递,它不可序列化。我必须通过一个变量传递它,然后它是可序列化的。这有点奇怪……但现在可以正常工作了。
标签: apache-spark serialization apache-kafka