Как сериализовать и десериализовать сообщение debezium с потоком QuarkusJAVA

Программисты JAVA общаются здесь
Anonymous
Как сериализовать и десериализовать сообщение debezium с потоком Quarkus

Сообщение Anonymous »

У меня возникла проблема с сериализацией и десериализацией с использованием потоков debezium и потоков Quarkus. Вот мой код:
1 — я написал свой pojo

Код: Выделить всё

@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
@RegisterForReflection
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonRootName("payload")
public class Payload {
public After after;
public Before before;
public Source source;
@JsonProperty("op")
public String op;

@Override
public String toString() {
return "Payload{" +
"after=" + after +
", before=" + before +
", source=" + source +
", op=" + op +
'}';
}
}

@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
@RegisterForReflection
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonRootName("after")
public class After {
@JsonProperty("id")
public int id;
@JsonProperty("redact_id")
public int redact_id;
@JsonProperty("descricao")
public String descricao;
@JsonProperty("data_inicial")
public String data_inicial;
@JsonProperty("data_final")
public String data_final;

@Override
public String toString() {
return "AfterFeriado{" +
"id=" + id +
", redact_id=" + redact_id +
", descricao=" + descricao +
", data_inicial=" + data_inicial +
", data_final=" + data_final +
'}';
}
}

@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
@RegisterForReflection
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonRootName("before")
public class Before {
@JsonProperty("id")
public String id;
@JsonProperty("redact_id")
public String redact_id;
@JsonProperty("descricao")
public String descricao;
@JsonProperty("data_inicial")
public String data_inicial;
@JsonProperty("data_final")
public String data_final;

@Override
public String toString() {
return "Before{" +
"id=" + id +
", redact_id=" + redact_id +
", descricao=" + descricao +
", data_inicial=" + data_inicial +
", data_final=" + data_final +
'}';
}
}

Я не знаю, нужно ли мне помещать @JsonProperty("id") в класс pojo.
Свойства приложения< /p>

Код: Выделить всё

quarkus.kafka-streams.bootstrap-servers=kafkabrokershere:9092
// others props

quarkus.kafka-streams.topics=collector_topic

и, наконец, моя функция для возврата KStream,

Код: Выделить всё

import java.time.LocalDateTime;
import java.util.Objects;

@ApplicationScoped
public class Topology {

private final String COLLECTOR = "collector_topic";
private final String UPDATE_TOPIC = "quarkus-udpated";
private final String REMOVED_TOPIC = "quarkus-removed";

@Produces
public org.apache.kafka.streams.Topology buildTopology() {
StreamsBuilder streamsBuilder = new StreamsBuilder();
ObjectMapperSerde payloadToProducer = new ObjectMapperSerde(PayloadToProducer.class);

var streamBuilder = this.buildStream(streamsBuilder);

System.out.println(streamBuilder.peek((k, v) -> System.out.println("hereee"+k)));
System.out.println(streamBuilder.peek((k, v) -> System.out.println("hereee"+v.op))); // null

if (streamBuilder == null) {
return null;
}
var operacaoDeCriacao = streamBuilder.filter((k,v) -> {
var created = Objects.equals(v.op, "c") || Objects.equals(v.op, "u");
System.out.println("null value here"+v);
return created;
});
var operacaoDeRemocao = streamBuilder.filter((k,v) ->  Objects.equals(v.op, "d"));
this.mapperStream(operacaoDeCriacao).to(UPDATE_TOPIC, Produced.with(Serdes.String(), payloadToProducer));
this.mapperStream(operacaoDeRemocao).to(REMOVED_TOPIC, Produced.with(Serdes.String(), payloadToProducer));

return streamsBuilder.build();
}

private KStream mapperStream(KStream stream) {
return stream
.map((k, v) -> {
// others operations to construct the new aggregated topic
var payloadToProducer = new PayloadToProducer(
2,2, 123, 123, "any valid information", LocalDateTime.now(), LocalDateTime.now()
);
return new KeyValue(k, payloadToProducer);
});
}

private KStream buildStream(StreamsBuilder streamsBuilder) {
// I try that
ObjectMapperSerde serdePayloadCurrent = new ObjectMapperSerde(Payload.class);

// and I try that
//Serde serdePayloadCurrent = Serdes.serdeFrom(new PayloadSerializer(), new PayloadDeserializer());

// And That
//JsonSerde payloadSerde = new JsonSerde(Payload.class);

return streamsBuilder
.stream(COLLECTOR, Consumed.with(Serdes.String(), serdePayloadCurrent))
.filter((k, v) -> {
System.out.println("value here null"+v.op);
return v != null;
});
}
}

сериализатор и десериализатор, я думаю, в этом проблема

Код: Выделить всё

public class PayloadDeserializer extends ObjectMapperDeserializer {
public PayloadDeserializer() {
super(Payload.class);
}
}

public class PayloadSerializer extends ObjectMapperSerializer {
}

Десериализатор и сериализатор должны
Отобразить serdePayloadCurrent = Serdes.serdeFrom(new PayloadSerializer(), new PayloadDeserializer()); и я поместил сюда, потому что пробовал раньше.
Затем возвращаю ноль в моем потоке, возможно, это проблема десериализатора

Подробнее здесь: https://stackoverflow.com/questions/787 ... kus-stream

Вернуться в «JAVA»