【问题标题】:Deserialization in Flink datastreamFlink 数据流中的反序列化
【发布时间】:2020-01-20 22:08:28
【问题描述】:

这里我给Kafka topic写了一个字符串,flink消费这个topic。反序列化是使用SimpleStringSchema完成的。当我需要消费整数值时,应该使用什么反序列化方法而不是SimpleStringschema???

DataStream<String> messageStream = env.addSource(new FlinkKafkaConsumer09<String>("test2", new SimpleStringSchema(), properties));

【问题讨论】:

  • 如何将 int 值写入 Kafka??
  • 使用kafka生产者代码:ProducerRecord record = new ProducerRecord(KafkaConstants.TOPIC_NAME, value);其中 value 是一个整数。这里我只写了 value 而不是 key。

标签: apache-kafka apache-flink


【解决方案1】:

您需要定义自己的 SerializationSchema

https://github.com/apache/flink/blob/master/flink-core/src/main/java/org/apache/flink/api/common/serialization/SerializationSchema.java#L32

或者将其保留为字符串,然后将流映射到您需要的类型

【讨论】:

  • 代码:DataStream messageStream = env.addSource(new FlinkKafkaConsumer09("test2", new DeserializationSchema(), properties));
  • @Noobie 你熟悉 Java 吗?你不能new 一个接口。如果您不至少 know the basics,您可能想跳过学习框架
猜你喜欢
  • 2020-08-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-11-18
  • 1970-01-01
  • 2013-08-20
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多