-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #17 from diegosneves/feat/kafka
Feat/kafka
- Loading branch information
Showing
13 changed files
with
301 additions
and
87 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1 @@ | ||
KAFKA_BOOTSTRAP_SERVERS=host.docker.internal:9094 |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,29 +1,20 @@ | ||
version: '3.9' | ||
|
||
services: | ||
email_server: | ||
image: diegoneves/email-server:latest | ||
restart: always | ||
container_name: email_server | ||
networks: | ||
- pedido-bridge | ||
ports: | ||
- "8081:8081" | ||
|
||
racha-pedido-app: | ||
image: diegoneves/racha-pedido:latest | ||
image: diegoneves/racha-pedido:beta | ||
container_name: racha-pedido | ||
ports: | ||
- "8080:8080" | ||
depends_on: | ||
- email_server | ||
networks: | ||
- pedido-bridge | ||
environment: | ||
- EMAIL_HOST=email_server | ||
- EMAIL_PORT=8081 | ||
|
||
|
||
networks: | ||
pedido-bridge: | ||
driver: bridge | ||
- KAFKA_BOOTSTRAP_SERVERS=host.docker.internal:9094 | ||
extra_hosts: | ||
- "host.docker.internal:172.17.0.1" |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
32 changes: 32 additions & 0 deletions
32
src/main/java/diegosneves/github/rachapedido/brokers/EmailSendConsumer.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,32 @@ | ||
package diegosneves.github.rachapedido.brokers; | ||
|
||
import diegosneves.github.rachapedido.model.NotificationEmail; | ||
import diegosneves.github.rachapedido.service.contract.EmailServiceContract; | ||
import diegosneves.github.rachapedido.utils.JsonAdapter; | ||
import org.apache.kafka.clients.consumer.ConsumerRecord; | ||
import org.springframework.beans.factory.annotation.Autowired; | ||
import org.springframework.kafka.annotation.EnableKafka; | ||
import org.springframework.kafka.annotation.KafkaListener; | ||
import org.springframework.stereotype.Component; | ||
|
||
@EnableKafka | ||
@Component | ||
public class EmailSendConsumer { | ||
|
||
|
||
private final EmailServiceContract emailService; | ||
private final JsonAdapter jsonAdapter; | ||
|
||
@Autowired | ||
public EmailSendConsumer(EmailServiceContract emailService, JsonAdapter jsonAdapter) { | ||
this.emailService = emailService; | ||
this.jsonAdapter = jsonAdapter; | ||
} | ||
|
||
@KafkaListener(topics = "${spring.kafka.topics.send_email}", groupId = "${spring.kafka.consumer.group-id}") | ||
public void sendEmail(ConsumerRecord<String, String> emailMessage) { | ||
NotificationEmail notification = jsonAdapter.fromJson(emailMessage.value(), NotificationEmail.class); | ||
this.emailService.sendEmail(notification); | ||
} | ||
|
||
} |
49 changes: 49 additions & 0 deletions
49
src/main/java/diegosneves/github/rachapedido/config/KafkaConfig.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,49 @@ | ||
package diegosneves.github.rachapedido.config; | ||
|
||
import org.apache.kafka.clients.consumer.ConsumerConfig; | ||
import org.apache.kafka.common.serialization.StringDeserializer; | ||
import org.springframework.beans.factory.annotation.Autowired; | ||
import org.springframework.beans.factory.annotation.Value; | ||
import org.springframework.context.annotation.Bean; | ||
import org.springframework.context.annotation.Configuration; | ||
import org.springframework.kafka.annotation.EnableKafka; | ||
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; | ||
import org.springframework.kafka.core.ConsumerFactory; | ||
import org.springframework.kafka.core.DefaultKafkaConsumerFactory; | ||
|
||
import java.util.HashMap; | ||
import java.util.Map; | ||
|
||
@Configuration | ||
@EnableKafka | ||
public class KafkaConfig { | ||
|
||
private final String bootstrapServers; | ||
private final String groupId; | ||
|
||
@Autowired | ||
public KafkaConfig(@Value("${spring.kafka.bootstrap-servers}") String bootstrapServers, @Value("${spring.kafka.consumer.group-id}") String groupId) { | ||
this.bootstrapServers = bootstrapServers; | ||
this.groupId = groupId; | ||
} | ||
|
||
@Bean | ||
public ConsumerFactory<String, String> consumerFactory() { | ||
Map<String, Object> props = new HashMap<>(); | ||
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.bootstrapServers); | ||
props.put(ConsumerConfig.GROUP_ID_CONFIG, this.groupId); | ||
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); | ||
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); | ||
return new DefaultKafkaConsumerFactory<>(props); | ||
} | ||
|
||
@Bean | ||
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { | ||
ConcurrentKafkaListenerContainerFactory<String, String> factory = | ||
new ConcurrentKafkaListenerContainerFactory<>(); | ||
factory.setConsumerFactory(consumerFactory()); | ||
factory.getContainerProperties().setMissingTopicsFatal(false); | ||
return factory; | ||
} | ||
|
||
} |
30 changes: 30 additions & 0 deletions
30
src/main/java/diegosneves/github/rachapedido/infrastructure/KafkaProducer.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,30 @@ | ||
package diegosneves.github.rachapedido.infrastructure; | ||
|
||
import diegosneves.github.rachapedido.utils.JsonAdapter; | ||
import org.springframework.beans.factory.annotation.Autowired; | ||
import org.springframework.kafka.core.KafkaTemplate; | ||
import org.springframework.stereotype.Component; | ||
|
||
@Component | ||
public class KafkaProducer { | ||
private final KafkaTemplate<String, String> kafkaTemplate; | ||
private final JsonAdapter jsonAdapter; | ||
|
||
@Autowired | ||
public KafkaProducer(KafkaTemplate<String, String> kafkaTemplate, JsonAdapter jsonAdapter) { | ||
this.kafkaTemplate = kafkaTemplate; | ||
this.jsonAdapter = jsonAdapter; | ||
} | ||
|
||
public void send(String topic, String message) { | ||
this.kafkaTemplate.send(topic, message); | ||
} | ||
|
||
public void convertObjectToJsonAndSend(String topic, Object object) { | ||
this.kafkaTemplate.send(topic, jsonAdapter.toJson(object)); | ||
} | ||
|
||
public void convertObjectToJsonAndSend(String topic, String key, Object object) { | ||
this.kafkaTemplate.send(topic, key, jsonAdapter.toJson(object)); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
33 changes: 33 additions & 0 deletions
33
src/main/java/diegosneves/github/rachapedido/utils/GsonAdapter.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,33 @@ | ||
package diegosneves.github.rachapedido.utils; | ||
|
||
import com.google.gson.Gson; | ||
import com.google.gson.GsonBuilder; | ||
import lombok.AllArgsConstructor; | ||
import lombok.Builder; | ||
import org.springframework.stereotype.Component; | ||
|
||
import java.time.LocalDateTime; | ||
|
||
@Component | ||
@Builder | ||
@AllArgsConstructor | ||
public class GsonAdapter implements JsonAdapter { | ||
|
||
private Gson gson; | ||
|
||
public GsonAdapter() { | ||
GsonBuilder gsonBuilder = new GsonBuilder(); | ||
gsonBuilder.registerTypeAdapter(LocalDateTime.class, new GsonLocalDateTime()); | ||
this.gson = gsonBuilder.setPrettyPrinting().create(); | ||
} | ||
|
||
@Override | ||
public String toJson(Object obj) { | ||
return this.gson.toJson(obj); | ||
} | ||
|
||
@Override | ||
public <T> T fromJson(String json, Class<T> classOfT) { | ||
return this.gson.fromJson(json, classOfT); | ||
} | ||
} |
28 changes: 28 additions & 0 deletions
28
src/main/java/diegosneves/github/rachapedido/utils/GsonLocalDateTime.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,28 @@ | ||
package diegosneves.github.rachapedido.utils; | ||
|
||
import com.google.gson.JsonDeserializationContext; | ||
import com.google.gson.JsonDeserializer; | ||
import com.google.gson.JsonElement; | ||
import com.google.gson.JsonParseException; | ||
import com.google.gson.JsonPrimitive; | ||
import com.google.gson.JsonSerializationContext; | ||
import com.google.gson.JsonSerializer; | ||
|
||
import java.lang.reflect.Type; | ||
import java.time.LocalDateTime; | ||
import java.time.format.DateTimeFormatter; | ||
|
||
public class GsonLocalDateTime implements JsonSerializer<LocalDateTime>, JsonDeserializer<LocalDateTime> { | ||
|
||
@Override | ||
public LocalDateTime deserialize(JsonElement jsonElement, Type type, JsonDeserializationContext jsonDeserializationContext) throws JsonParseException { | ||
String ldtString = jsonElement.getAsString(); | ||
return LocalDateTime.parse(ldtString, DateTimeFormatter.ISO_LOCAL_DATE_TIME); | ||
} | ||
|
||
@Override | ||
public JsonElement serialize(LocalDateTime localDateTime, Type type, JsonSerializationContext jsonSerializationContext) { | ||
return new JsonPrimitive(localDateTime.format(DateTimeFormatter.ISO_LOCAL_DATE_TIME)); | ||
} | ||
|
||
} |
9 changes: 9 additions & 0 deletions
9
src/main/java/diegosneves/github/rachapedido/utils/JsonAdapter.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,9 @@ | ||
package diegosneves.github.rachapedido.utils; | ||
|
||
public interface JsonAdapter { | ||
|
||
String toJson(Object obj); | ||
|
||
<T> T fromJson(String json, Class<T> classOfT); | ||
|
||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.