【问题标题】:How to Serialize String in Scala如何在 Scala 中序列化字符串
【发布时间】: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
  }

Error publishing in kafka

【问题讨论】:

  • 你为什么不分享你制作人的代码?
  • Spark Streaming 已弃用。你为什么不使用结构化流?并且请不要将图片用于错误
  • 很抱歉没有提供有关我的错误的更多详细信息。我从老师那里收到了关于问题所在的答复。问题在于火花上下文。我将 spark 上下文作为参数传递,它不可序列化。我必须通过一个变量传递它,然后它是可序列化的。这有点奇怪……但现在可以正常工作了。

标签: apache-spark serialization apache-kafka


【解决方案1】:

错误并没有说明字符串。一个 Dstream 是不可序列化的,这意味着你已经在 executor 范围之间共享了它的一部分,这让 Spark 认为它需要序列化它

您确实应该显示所有代码,但是要使用 KafkaProducer 是 Spark Streaming,您需要使用 foreachPartition,然后在该块内创建一个 Producer。

对于每个分区,循环遍历每个 RDD,然后使用 KafkaProducer.send 方法

除非你想定义自己的序列化,否则你不需要担心序列化

【讨论】:

    【解决方案2】:

    如果没有代码,我无法说出确切的解决方案。我认为您的问题与 Kafka 的属性无关。

    在错误日志中,Spark 尝试序列化一个类但失败了。

    请检查您在 foreachRDD 块中的代码。我认为您使用了不可序列化的类。如果可以,请检查您的类并将 Serializable 实现添加到您的类中。或者尝试使用 String 类型。

    【讨论】:

    • 最终问题出在火花上下文上。我将 spark 上下文作为参数传递,而 spark 上下文不可序列化。
    猜你喜欢
    • 1970-01-01
    • 2020-09-27
    • 1970-01-01
    • 2017-10-09
    • 2018-10-22
    • 2010-11-22
    • 2018-03-17
    • 1970-01-01
    • 2020-09-14
    相关资源
    最近更新 更多