Flink Push Row to kafka

Viewed 164

I have a flink Row and column names, This Row can be accessed via field names or index. I want to sink it into kafka in JSON using vanilla flink kafka producer. How can I do it? Is destination Json Schema require it to sink it to kafka?

1 Answers

You would need to provide schema for the kafka producer you specify for the stream. Fortunately, Flink does provide you with schemas that you can modify as required. If you just want to send your Object as a JSON string, you can do something like this:

Convert your stream of Objects to a stream of String JSONs as such:

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);
    }
});

You can then define a Kafka Producer as follows:

    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);
    }

Finally specify the kafka producer to the stream and use it to sink your JSON string messages:

jsons
    .addSink(FlinkUtils.getKafkaProducer("outputTopic))
    .name("JSON messages sink");
Related