【问题标题】:Flink Streaming: Unexpected charaters in serialized String messagesFlink Streaming:序列化字符串消息中的意外字符
【发布时间】:2018-02-04 22:31:59
【问题描述】:

我的流正在生成Tuple2<String,String> 类型的记录

.toString() 输出(usr12345,{"_key":"usr12345","_temperature":46.6})

其中键是usr12345,值是{"_key":"usr12345","_temperature":46.6}

流上的.print()正确输出值:

(usr12345,{"_key":"usr12345","_temperature":46.6})

但是当我将流写入 Kafka 时,键变为 usr12345(开头有一个空格)和值 ({"_key":"usr12345","_temperature":46.6}

注意键开头的空格和值开头的左括号。

很奇怪。为什么会发生这种情况?

序列化代码如下:

TypeInformation<String> resultType = TypeInformation.of(String.class);

KeyedSerializationSchema<Tuple2<String, String>> schema =
      new TypeInformationKeyValueSerializationSchema<>(resultType, resultType, env.getConfig());

FlinkKafkaProducer010.FlinkKafkaProducer010Configuration flinkKafkaProducerConfig = FlinkKafkaProducer010.writeToKafkaWithTimestamps(
      stream,   
      "topic",    
      schema,  
      kafkaProducerProperties);

【问题讨论】:

  • 你描述的有点奇怪,你有没有尝试过创建一个kafka sink并做stream.addsink(kafkaSink)?能不能解决问题?
  • @BiplobBiswas 好吧,我按照 Flink Kafka 文档中描述的说明进行操作。 ci.apache.org/projects/flink/flink-docs-release-1.3/dev/…据此,这是使用我正在使用的 Java Kafka 0.10+ 的正确方法。

标签: serialization apache-kafka apache-flink kafka-producer-api flink-streaming


【解决方案1】:

TypeInformationKeyValueSerializationSchema 使用 Flink 的自定义序列化器序列化数据,这意味着结果必须被解释为二进制数据。 Flink 的 String 序列化器写入 String 的长度,然后编码所有字符。

我假设您使用纯字符串反序列化器反序列化 Kafka 主题。对于键,序列化长度被解释为空白字符。对于值,长度被解释为'('

尝试使用不同的序列化器将键和值序列化为纯字符串或使用兼容的反序列化器。

【讨论】:

    猜你喜欢
    • 2017-09-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多