【问题标题】:is JSONDeserializationSchema() deprecated in Flink?Flink 中不推荐使用 JSONDeserializationSchema() 吗?
【发布时间】:2020-05-28 14:54:36
【问题描述】:

我是 Flink 新手,做的事情与下面的链接非常相​​似。

Cannot see message while sinking kafka stream and cannot see print message in flink 1.2

我还尝试添加 JSONDeserializationSchema() 作为我的 Kafka 输入 JSON 消息的反序列化程序,该消息没有密钥。

但我发现 JSONDeserializationSchema() 不存在。

如果我做错了什么,请告诉我。

【问题讨论】:

  • Flink 中没有JSONDeserializationSchema 这样的东西。有JSONKeyValueDeserializationSchema,你是这个意思吗?
  • 通过这个 stackoverflow ans stackoverflow.com/questions/45564373/… 我正在尝试使用 JSONKeyValueDeserializationSchema 不起作用的非密钥 json 消息。
  • 另外,GitHub 上的代码似乎正在使用 JSONDeserializationSchema github.com/masato/streams-flink-scala-examples/blob/master/src/…
  • 啊,是的,对不起,我的错。我正在查看 Flink 1.8,但正如 @David 下面提到的,它已经从这个版本中删除了。
  • 不是问题@DominikWosiński :)

标签: apache-flink flink-streaming


【解决方案1】:

JSONDeserializationSchema 在 Flink 1.8 中被删除,此前已被弃用。

推荐的方法是编写一个实现DeserializationSchema<T> 的反序列化器。这是我从Flink Operations Playground 复制的示例:

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;

import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;

import java.io.IOException;

/**
 * A Kafka {@link DeserializationSchema} to deserialize {@link ClickEvent}s from JSON.
 *
 */
public class ClickEventDeserializationSchema implements DeserializationSchema<ClickEvent> {

    private static final long serialVersionUID = 1L;

    private static final ObjectMapper objectMapper = new ObjectMapper();

    @Override
    public ClickEvent deserialize(byte[] message) throws IOException {
        return objectMapper.readValue(message, ClickEvent.class);
    }

    @Override
    public boolean isEndOfStream(ClickEvent nextElement) {
        return false;
    }

    @Override
    public TypeInformation<ClickEvent> getProducedType() {
        return TypeInformation.of(ClickEvent.class);
    }
}

对于 Kafka 生产者,您需要实现 KafkaSerializationSchema&lt;T&gt;,您会在同一个项目中找到这样的示例。

【讨论】:

  • 谢谢大卫?。我们有什么其他的地方或任何解决方法吗?我正在尝试使用非密钥 json 消息。
  • 我已更新我的答案以包含一个示例。这种方法应该有更好的性能,并且更干净地与 Kafka 消费者集成。
  • 谢谢老兄!!这对我非常有帮助,解决了我几个小时都在努力解决的问题。干杯
【解决方案2】:

为了解决从 Kafka 读取非关键 JSON 消息的问题,我使用了案例类和 JSON 解析器。

以下代码创建一个案例类并使用 play API 解析 JSON 字段。

import play.api.libs.json.JsValue

object CustomerModel {

  def readElement(jsonElement: JsValue): Customer = {
    val id = (jsonElement \ "id").get.toString().toInt
    val name = (jsonElement \ "name").get.toString()
    Customer(id,name)
  }
case class Customer(id: Int, name: String)
}

def main(args: Array[String]): Unit = {
val env = StreamExecutionEnvironment.getExecutionEnvironment
val properties = new Properties()
properties.setProperty("bootstrap.servers", "xxx.xxx.0.114:9092")
properties.setProperty("group.id", "test-grp")

val consumer = new FlinkKafkaConsumer[String]("customer", new SimpleStringSchema(), properties)
val stream1 = env.addSource(consumer).rebalance

val stream2:DataStream[Customer]= stream1.map( str =>{Try(CustomerModel.readElement(Json.parse(str))).getOrElse(Customer(0,Try(CustomerModel.readElement(Json.parse(str))).toString))
    })

stream2.print("stream2")
env.execute("This is Kafka+Flink")

}

Try 方法可以让您克服解析数据时抛出的异常 并在其中一个字段中返回异常(如果我们需要),否则它可以只返回带有任何给定或默认字段的案例类对象。

代码的示例输出为:

stream2:1> Customer(1,"Thanh")
stream2:1> Customer(5,"Huy")
stream2:3> Customer(0,Failure(com.fasterxml.jackson.databind.JsonMappingException: No content to map due to end-of-input
 at [Source: ; line: 1, column: 0]))

我不确定这是否是最好的方法,但它现在对我有用。

【讨论】:

    猜你喜欢
    • 2011-01-25
    • 1970-01-01
    • 1970-01-01
    • 2021-05-22
    • 2021-01-03
    • 1970-01-01
    • 1970-01-01
    • 2017-04-05
    • 1970-01-01
    相关资源
    最近更新 更多