【发布时间】:2023-03-11 09:10:02
【问题描述】:
我有一个 flink Row 和列名,这个Row 可以通过字段名或索引访问。我想使用 vanilla flink kafka producer 将它放入 JSON 中的 kafka 中。我该怎么做?目标 Json Schema 是否需要将其下沉到 kafka?
【问题讨论】:
标签: apache-flink flink-streaming
我有一个 flink Row 和列名,这个Row 可以通过字段名或索引访问。我想使用 vanilla flink kafka producer 将它放入 JSON 中的 kafka 中。我该怎么做?目标 Json Schema 是否需要将其下沉到 kafka?
【问题讨论】:
标签: apache-flink flink-streaming
您需要为您为流指定的 kafka 生产者提供架构。幸运的是,Flink 确实为您提供了可以根据需要进行修改的模式。如果您只想将您的对象作为 JSON 字符串发送,您可以执行以下操作:
将您的对象流转换为字符串 JSON 流,如下所示:
SingleOutputStreamOperator<String> jsons = dataStream.map(new MapFunction<Object, String>() {
@Override
public String map(Object value) throws Exception {
// Gson creation can be put in a static utility method if you want to avoid recreating
Gson gson = new GsonBuilder().create();
return gson.toJson(value);
}
});
然后您可以按如下方式定义 Kafka Producer:
public static FlinkKafkaProducer<String> getKafkaProducer(String topic) {
String kafkaBootstrapServers = "localhost:9092";
String kafkaGroup = "kafkaGroup";
Properties propertiesProducer = new Properties();
propertiesProducer.setProperty("bootstrap.servers", kafkaBootstrapServers);
propertiesProducer.setProperty("group.id", kafkaGroup);
SimpleStringSchema simpleStringSchema = new SimpleStringSchema() {
public String deserialize(byte[] message) {
return message == null ? null : super.deserialize(message);
}
};
return new FlinkKafkaProducer(topic, simpleStringSchema, propertiesProducer);
}
最后为流指定 kafka 生产者并使用它来接收您的 JSON 字符串消息:
jsons
.addSink(FlinkUtils.getKafkaProducer("outputTopic))
.name("JSON messages sink");
【讨论】: