我编写了一个简单的卡夫卡序列化程序,以检查序列化程序-->configure方法可以告诉我当前的键是被序列化还是被赋值。
看看接口定义,应该是这样的,但是当我使用下面的代码时,isKey总是返回false。
有人能告诉我什么时候会触发配置方法吗?
我是否误解了isKey变量,它实际上表明了一些不同的东西?
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import org.apache.kafka.common.errors.SerializationException;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.connect.data.Struct;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.util.Map;
public class KafkaSerializer<T> implements Serializer<T> {
protected ObjectMapper objectMapper = new ObjectMapper();
protected JsonSchemaUtils jsonSchema = new JsonSchemaUtils();
private boolean isKey;
@Override
public void configure(Map<String, ?> configs, boolean isKey) {
this.isKey = isKey;
Serializer.super.configure(configs, isKey);
}
@Override
public byte[] serialize(String topic, T data) {
Struct value = (Struct) data;
System.out.println("isKey: " + isKey);
ObjectNode schemaNode = jsonSchema.envelopeSchema(value);
ObjectNode payloadNode = jsonSchema.envelopePayload(value);
ByteArrayOutputStream out = new ByteArrayOutputStream();
ObjectNode result = jsonSchema.envelope(schemaNode, payloadNode);
try {
out.write(objectMapper.writeValueAsBytes(result));
byte[] bytes = out.toByteArray();
out.close();
return bytes;
} catch (JsonProcessingException e) {
throw new SerializationException(e);
} catch (IOException e) {
e.printStackTrace();
}
return new byte[0];
}
@Override
public byte[] serialize(String topic, Headers headers, T data) {
return Serializer.super.serialize(topic, headers, data);
}
@Override
public void close() {
Serializer.super.close();
}
}