mirror of
https://github.com/Example-uPagge/quarkus-kafka-panache.git
synced 2024-06-14 11:52:55 +03:00
solution
This commit is contained in:
parent
52ab41e68e
commit
84a83cd574
@ -21,7 +21,7 @@ public class KafkaMessageGenerator {
|
||||
|
||||
|
||||
public void start(@Observes StartupEvent ev) {
|
||||
for (int i = 0; i < 20; i++) {
|
||||
for (int i = 0; i < 50; i++) {
|
||||
final KafkaMessage kafkaMessage = new KafkaMessage();
|
||||
kafkaMessage.setCount(i);
|
||||
telegramChannel.sendAndAwait(kafkaMessage);
|
||||
|
@ -0,0 +1,21 @@
|
||||
package dev.struchkov.example.quarkus.kafka;
|
||||
|
||||
import io.smallrye.mutiny.Uni;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.hibernate.reactive.mutiny.Mutiny;
|
||||
|
||||
import javax.enterprise.context.ApplicationScoped;
|
||||
|
||||
@ApplicationScoped
|
||||
@RequiredArgsConstructor
|
||||
public class EntityRepositoryImpl {
|
||||
|
||||
private final Mutiny.SessionFactory factory;
|
||||
|
||||
public Uni<EntityForDb> save(EntityForDb entityForDb) {
|
||||
return factory.withTransaction(
|
||||
session -> session.merge(entityForDb)
|
||||
);
|
||||
}
|
||||
|
||||
}
|
@ -10,14 +10,14 @@ import javax.enterprise.context.ApplicationScoped;
|
||||
@RequiredArgsConstructor
|
||||
public class KafkaHandler {
|
||||
|
||||
private final PanacheRepositoryImpl panacheRepository;
|
||||
private final EntityRepositoryImpl panacheRepository;
|
||||
|
||||
@Incoming("test")
|
||||
public Uni<Void> handle(KafkaMessage message) {
|
||||
System.out.println("Получено сообщение " + message);
|
||||
final EntityForDb entityForDb = new EntityForDb();
|
||||
entityForDb.setCount(message.getCount());
|
||||
return panacheRepository.persistAndFlush(entityForDb).replaceWithVoid();
|
||||
return panacheRepository.save(entityForDb).replaceWithVoid();
|
||||
}
|
||||
|
||||
}
|
||||
|
@ -1,10 +0,0 @@
|
||||
package dev.struchkov.example.quarkus.kafka;
|
||||
|
||||
import io.quarkus.hibernate.reactive.panache.PanacheRepository;
|
||||
|
||||
import javax.enterprise.context.ApplicationScoped;
|
||||
|
||||
@ApplicationScoped
|
||||
public class PanacheRepositoryImpl implements PanacheRepository<EntityForDb> {
|
||||
|
||||
}
|
@ -6,6 +6,7 @@ quarkus:
|
||||
jdbc: false
|
||||
reactive:
|
||||
url: postgresql://localhost:5432/quarkus_test
|
||||
max-size: 20
|
||||
hibernate-orm:
|
||||
database:
|
||||
generation: update
|
||||
|
Loading…
Reference in New Issue
Block a user