【发布时间】:2018-12-10 21:14:55
【问题描述】:
自定义类
人物
class Person
{
private Integer id;
private String name;
//getters and setters
}
Kafka Flink 连接器
TypeInformation<Person> info = TypeInformation.of(Person.class);
TypeInformationSerializationSchema schema = new TypeInformationSerializationSchema(info, new ExecutionConfig());
DataStream<Person> input = env.addSource( new FlinkKafkaConsumer08<>("persons", schema , getKafkaProperties()));
现在如果我发送下面的 json
{ "id" : 1, "name": Synd }
通过Kafka Console Producer,flink代码抛出空指针异常
但是,如果我使用 SimpleStringSchema 而不是之前定义的 CustomSchema,则会打印流。
上面的设置有什么问题
【问题讨论】:
标签: json serialization apache-kafka deserialization apache-flink