diff --git a/clearing-parent/balance-service/pom.xml b/clearing-parent/balance-service/pom.xml
index 4414833e5..5b4cd23de 100644
--- a/clearing-parent/balance-service/pom.xml
+++ b/clearing-parent/balance-service/pom.xml
@@ -14,6 +14,7 @@
17
17
+ 3.0.1
@@ -56,8 +57,24 @@
assertj-core
test
-
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+ org.springframework.kafka
+ spring-kafka-test
+ 2.8.8
+ test
+
+
+ org.springframework.kafka
+ spring-kafka
+ 2.8.8
+
+
diff --git a/clearing-parent/balance-service/src/main/resources/application.properties b/clearing-parent/balance-service/src/main/resources/application.properties
index 46ac587bd..3751e0dbb 100644
--- a/clearing-parent/balance-service/src/main/resources/application.properties
+++ b/clearing-parent/balance-service/src/main/resources/application.properties
@@ -2,7 +2,6 @@ spring.main.web-application-type=none
balance-service.hazelcast.cluster-members=127.0.0.1:5701
balance-service.hazelcast.login=dev
balance-service.hazelcast.password=dev-pass
-
balance-service.kafka-consumer.bootstrap-servers=localhost:9092
balance-service.kafka-consumer.group-id=dev-group-balance-service
balance-service.kafka-consumer.enable-auto-commit=false
@@ -10,10 +9,11 @@ balance-service.kafka-consumer.session-timeout-ms=30000
balance-service.kafka-consumer.auto-offset-reset=latest
balance-service.kafka-consumer.linger-ms=1
balance-service.kafka-consumer.buffer-memory=33554432
-
balance-service.kafka-producer.bootstrap-servers=localhost:9092
balance-service.kafka-producer.acks=all
balance-service.kafka-producer.retries=0
balance-service.kafka-producer.batch-size=16384
balance-service.kafka-producer.linger-ms=1
-balance-service.kafka-producer.buffer-memory=33554432
\ No newline at end of file
+balance-service.kafka-producer.buffer-memory=33554432
+#for testing
+spring.kafka.consumer.group-id=dev-group-balance-service
\ No newline at end of file
diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java
index b56fef2b1..256267d20 100644
--- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java
+++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java
@@ -1,24 +1,17 @@
package ru.spcex.clearing.balance.config;
-import org.apache.kafka.clients.consumer.MockConsumer;
-import org.apache.kafka.clients.consumer.OffsetResetStrategy;
-import org.apache.kafka.clients.producer.MockProducer;
-import org.apache.kafka.clients.producer.Producer;
-import org.apache.kafka.common.serialization.StringSerializer;
-import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
-import ru.spcex.clearing.platform.messaging.serialization.JsonSerializer;
@Configuration
public class KafkaTestConfig {
- @Bean
- public MockConsumer createTestConsumer() {
- return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
- }
-
- @Bean
- public Producer createTestProducer() {
- return new MockProducer<>(true, new StringSerializer(), new JsonSerializer());
- }
+// @Bean
+// public MockConsumer createTestConsumer() {
+// return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
+// }
+//
+// @Bean
+// public Producer createTestProducer() {
+// return new MockProducer<>(true, new StringSerializer(), new JsonSerializer());
+// }
}
diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConsumer.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConsumer.java
new file mode 100644
index 000000000..502e8a50a
--- /dev/null
+++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConsumer.java
@@ -0,0 +1,41 @@
+package ru.spcex.clearing.balance.config;
+
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.kafka.annotation.KafkaListener;
+import org.springframework.stereotype.Component;
+
+import java.util.concurrent.CountDownLatch;
+
+@Component
+public class KafkaTestConsumer {
+
+ public static final String TOPIC_GALB = "launcher-GALB";
+ private static final Logger LOGGER = LoggerFactory.getLogger(KafkaTestConsumer.class);
+ private final String TOPIC_NAME = "com.madadipouya.kafka.user";
+ private CountDownLatch latch = new CountDownLatch(1);
+
+ private String payload;
+
+ @KafkaListener(topics = TOPIC_GALB)
+ public void receive(ConsumerRecord, ?> consumerRecord) {
+ LOGGER.info("received payload='{}'", consumerRecord.toString());
+
+ payload = consumerRecord.toString();
+ latch.countDown();
+ }
+
+ public CountDownLatch getLatch() {
+ return latch;
+ }
+
+ public void resetLatch() {
+ latch = new CountDownLatch(1);
+ }
+
+ public String getPayload() {
+ return payload;
+ }
+
+}
diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/EmbeddedKafkaTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/EmbeddedKafkaTest.java
new file mode 100644
index 000000000..8c5b3dd79
--- /dev/null
+++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/EmbeddedKafkaTest.java
@@ -0,0 +1,104 @@
+package ru.spcex.clearing.balance.service;
+
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.producer.Producer;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.serialization.StringDeserializer;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.junit.jupiter.api.*;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
+import org.springframework.kafka.core.DefaultKafkaProducerFactory;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.kafka.listener.ContainerProperties;
+import org.springframework.kafka.listener.KafkaMessageListenerContainer;
+import org.springframework.kafka.listener.MessageListener;
+import org.springframework.kafka.test.EmbeddedKafkaBroker;
+import org.springframework.kafka.test.context.EmbeddedKafka;
+import org.springframework.kafka.test.utils.ContainerTestUtils;
+import org.springframework.kafka.test.utils.KafkaTestUtils;
+import org.springframework.test.context.junit.jupiter.SpringExtension;
+import ru.spcex.clearing.balance.config.KafkaTestConsumer;
+
+import java.io.File;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+
+import static org.hamcrest.CoreMatchers.containsString;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static ru.spcex.clearing.balance.config.KafkaTestConsumer.TOPIC_GALB;
+
+@EmbeddedKafka//(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})
+@SpringBootTest(properties = "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}")
+@ExtendWith(SpringExtension.class)
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+class EmbeddedKafkaTest {
+ static {
+ Path path = Paths.get("src", "main", "resources");
+ String currentPath = path.toAbsolutePath().toString();
+ System.setProperty("spring.config.location", currentPath + File.separator);
+ }
+
+ @Autowired
+ public KafkaTemplate template;
+ BlockingQueue> records;
+ KafkaMessageListenerContainer container;
+ @Autowired
+ private EmbeddedKafkaBroker embeddedKafkaBroker;
+ @Autowired
+ private KafkaTestConsumer consumer;
+
+ @BeforeEach
+ void setup() {
+ consumer.resetLatch();
+ }
+
+ @BeforeAll
+ void setUp() {
+ Map configs = new HashMap<>(KafkaTestUtils.consumerProps("consumer", "false", embeddedKafkaBroker));
+ DefaultKafkaConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>(configs, new StringDeserializer(), new StringDeserializer());
+ ContainerProperties containerProperties = new ContainerProperties(TOPIC_GALB);
+ container = new KafkaMessageListenerContainer<>(consumerFactory, containerProperties);
+ records = new LinkedBlockingQueue<>();
+ container.setupMessageListener((MessageListener) records::add);
+ container.start();
+ ContainerTestUtils.waitForAssignment(container, embeddedKafkaBroker.getPartitionsPerTopic());
+ }
+
+ @AfterAll
+ void tearDown() {
+ container.stop();
+ }
+
+ @Test
+ public void kafkaSetup_withTopic_ensureSendMessageIsReceived() throws Exception {
+ // Arrange
+ Map configs = new HashMap<>(KafkaTestUtils.producerProps(embeddedKafkaBroker));
+ Producer producer = new DefaultKafkaProducerFactory<>(configs, new StringSerializer(), new StringSerializer()).createProducer();
+
+ String data = "Sending with default template";
+ //Act
+ producer.send(new ProducerRecord<>(TOPIC_GALB, "my-aggregate-id", data));
+// producer.flush();
+
+// template.send(TOPIC, data);
+
+ // Assert
+ ConsumerRecord singleRecord = records.poll(10, TimeUnit.SECONDS);
+ boolean messageConsumed = consumer.getLatch()
+ .await(30, TimeUnit.SECONDS);
+ assertTrue(messageConsumed);
+ assertThat(consumer.getPayload(), containsString(data));
+// assertThat(singleRecord).isNotNull();
+// assertThat(singleRecord.key()).isEqualTo("my-aggregate-id");
+// assertThat(singleRecord.value()).isEqualTo("{\"event\":\"Test Event\"}");
+ }
+}
\ No newline at end of file
diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceKafkaTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceKafkaTest.java
new file mode 100644
index 000000000..bc17da0da
--- /dev/null
+++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceKafkaTest.java
@@ -0,0 +1,69 @@
+package ru.spcex.clearing.balance.service;
+
+import org.apache.kafka.clients.producer.Producer;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.kafka.core.DefaultKafkaProducerFactory;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.kafka.test.EmbeddedKafkaBroker;
+import org.springframework.kafka.test.context.EmbeddedKafka;
+import org.springframework.kafka.test.utils.KafkaTestUtils;
+import org.springframework.test.context.junit.jupiter.SpringExtension;
+import ru.spcex.clearing.balance.config.KafkaTestConsumer;
+
+import java.io.File;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
+import static org.hamcrest.CoreMatchers.containsString;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static ru.spcex.clearing.balance.config.KafkaTestConsumer.TOPIC_GALB;
+
+@EmbeddedKafka//(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})
+@SpringBootTest(properties = "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}")
+@ExtendWith(SpringExtension.class)
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+class Sdf08ServiceKafkaTest {
+ static {
+ Path path = Paths.get("src", "main", "resources");
+ String currentPath = path.toAbsolutePath().toString();
+ System.setProperty("spring.config.location", currentPath + File.separator);
+ }
+
+ @Autowired
+ public KafkaTemplate template;
+
+ @Autowired
+ private EmbeddedKafkaBroker embeddedKafkaBroker;
+ @Autowired
+ private KafkaTestConsumer consumer;
+
+ @Test
+ public void newSDf08() throws Exception {
+ // Arrange
+ Map configs = new HashMap<>(KafkaTestUtils.producerProps(embeddedKafkaBroker));
+ Producer producer = new DefaultKafkaProducerFactory<>(configs, new StringSerializer(), new StringSerializer()).createProducer();
+
+
+ String data = "Sending with default template";
+ //Act
+ producer.send(new ProducerRecord<>(TOPIC_GALB, "my-aggregate-id", data));
+// producer.flush();
+// template.send(TOPIC, data);
+
+ // Assert
+ boolean messageConsumed = consumer.getLatch()
+ .await(10, TimeUnit.SECONDS);
+ assertTrue(messageConsumed);
+ assertThat(consumer.getPayload(), containsString(data));
+ }
+}
\ No newline at end of file