package slite.lib_java;

import com.google.gson.JsonElement;
import com.google.gson.JsonParser;
import com.google.gson.Gson;

public abstract class SliteKafkaClientJson extends SliteKafkaClient<String, String>
{
	private Gson gson = new Gson();

	public abstract void receive(String topic, String key, long timestampMillis, JsonElement data) throws Exception;

	public SliteKafkaClientJson() throws Exception
	{
		this.keySerializer = "String";
		this.valueSerializer = "String";
		this.init();
	}

	@Override
	public void receive(String topic, String key, long timestampMillis, String data) throws Exception
	{
		JsonElement obj = JsonParser.parseString(data);
		this.receive(topic, key, timestampMillis, obj);
	}

	public void produce(JsonElement data, String key, String topic, Long timestamp) throws Exception
	{
		this.produce(this.gson.toJson(data), key, topic, timestamp);
	}

	public void produce(JsonElement data, String key, String topic) throws Exception
	{
		this.produce(this.gson.toJson(data), key, topic);
	}

	public void produce(JsonElement data, String key) throws Exception
	{
		this.produce(this.gson.toJson(data), key);
	}
}