Как соединить два Kafka topic через KTable?
Есть две темы в Kafka? которые пишутcя с использованием Avro схем (для key и value) Эти топики нужно объединить в один.
Как я понимаю, это нужно сделать через KTable и leftJoin. Из того, что мне удалось понять, из вычитанного на просторах интернета, я делаю следующее:
public static Topology createTopology() {
StreamsBuilder builder = new StreamsBuilder();
final Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url","http://localhost:8081");
final Serde<GenericRecord> keyGenericAvroSerde = new GenericAvroSerde();
keyGenericAvroSerde.configure(serdeConfig, true);
final Serde<GenericRecord> valueGenericAvroSerde = new GenericAvroSerde();
valueGenericAvroSerde.configure(serdeConfig, false);
KeyValueBytesStoreSupplier storeSupplierInfoReg = Stores.persistentKeyValueStore("top_1");
KeyValueBytesStoreSupplier storeSupplierDoc = Stores.persistentKeyValueStore("top_2");
KTable<GenericRecord, GenericRecord> leftTable = builder.table("top_1",
Materialized.<GenericRecord, GenericRecord>as(storeSupplierInfoReg).withKeySerde(keyGenericAvroSerde).withValueSerde(valueGenericAvroSerde));
KTable<GenericRecord, GenericRecord> rightTable = builder.table("top_2",
Materialized.<GenericRecord, GenericRecord>as(storeSupplierDoc).withKeySerde(keyGenericAvroSerde).withValueSerde(valueGenericAvroSerde));
Schema schema = createMergeSchema();
KTable<GenericRecord, Object> joined = leftTable.leftJoin(rightTable, (o1, o2) -> {
final GenericRecord viewRegion = new GenericData.Record(schema);
viewRegion.put("left", o1);
viewRegion.put("right", o2);
return viewRegion;
});
joined.toStream().to("result");
return builder.build();
}
static Properties getProperties() {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG,"test1");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");
props.put(KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, GenericAvroSerde.class);
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, GenericAvroSerde.class);
return props; }
получаю:
org.apache.kafka.common.errors.SerializationException: Error registering Avro schema: {"type":"record","name":"top_1Key","namespace":"W_","fields":[{"name":"Ref","type":"string","doc":"Document.Services"}]}
Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Schema being registered is incompatible with an earlier schema; error code: 409
Собственно вопросы: как сделать join двух KTable? У которых ключи и значения не примитивные типы, а объекты построенные по Avro схеме?