string in double quotes after a map operation in Kafka streams

Viewed 2011

I use Kafka streams DSL and map to convert a KStream<String, JsonNode> to KStream<String, String>.

In ValueMapper function, I simple return new ValueMapper("key", "some constant string"), but where the value is sent back to Kafka use KStream.to("some topic"), I got the result is being added double quotes.

my code is like this:

KStream<String, JsonNode> views = builder.stream("fromTopic");
views.map(new ValueKMapper()).to("toTopic");

and ValueMapper simply implements KeyValueMapper and the code of apply() is:

public KeyValue<String, String> apply(String key, JsonNode value) 
{
    return new KeyValue<String,String>("a", "hello");
}

and then when I consume the toTopic, I got the ""hello"", with quotes added.

Maybe it's a bug of Kafka streams?

2 Answers

I assume that the method apply() you've provided in your question

public KeyValue<String, String> apply(String key, JsonNode value) 
{
    return new KeyValue<String,String>("a", "hello");
}

is not passing hardcoded key and value to KeyValue's constructor. My guess is that the problem has something to do with JsonNode. Maybe, the actual implementation of your method uses value.get(key) i.e.

public KeyValue<String, String> apply(String key, JsonNode value) 
{
    return new KeyValue<String,String>(key, value.get(key));
}

However, value.get(key) will return a TextNode and the toString() method will return a string representation of TextNode including quotes. In order to parse JsonNode properly, you need to use textValue() method so your method will become

public KeyValue<String, String> apply(String key, JsonNode value) 
{
    return new KeyValue<String,String>(key, value.get(key).textValue());
}

Example: Assuming you have a key a and a value hello,

json.get("a").toString())

will return "hello" while

json.get("a").textValue();

will return hello

It could happen that you are using Json serde as default. Try adding string serde when writing to sink:

.to("toTopic", Produced.with(Serdes.String(), Serdes.String()));
Related