【发布时间】: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