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 +
'}';
}
}
Свойства приложения< /p>
Код: Выделить всё
quarkus.kafka-streams.bootstrap-servers=kafkabrokershere:9092
// others props
quarkus.kafka-streams.topics=collector_topic
Код: Выделить всё
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