simple cheng.
This commit is contained in:
parent
0de9b4cc5c
commit
b58ed3023d
2 changed files with 11 additions and 17 deletions
|
|
@ -1,6 +1,5 @@
|
||||||
package ru.spcex.clearing.balance.config;
|
package ru.spcex.clearing.balance.config;
|
||||||
|
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
|
||||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||||
import org.apache.kafka.clients.producer.MockProducer;
|
import org.apache.kafka.clients.producer.MockProducer;
|
||||||
|
|
@ -17,7 +16,7 @@ public class KafkaTestConfig {
|
||||||
|
|
||||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||||
@Bean
|
@Bean
|
||||||
public Consumer<String, Object> createTestConsumer() {
|
public MockConsumer<String, Object> createTestConsumer() {
|
||||||
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,18 +1,16 @@
|
||||||
package ru.spcex.clearing.balance.service;
|
package ru.spcex.clearing.balance.service;
|
||||||
|
|
||||||
import com.hazelcast.core.IMap;
|
import com.hazelcast.core.IMap;
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||||
|
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||||
import org.apache.kafka.common.TopicPartition;
|
import org.apache.kafka.common.TopicPartition;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import ru.clearing.classes.statics.data.account.BankAccount;
|
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf08;
|
import ru.clearing.classes.statics.data.sdf.SDf08;
|
||||||
import ru.clearing.classes.statics.data.statement.Statement;
|
import ru.clearing.classes.statics.data.statement.Statement;
|
||||||
import ru.spcex.clearing.balance.utils.ImapEvent;
|
import ru.spcex.clearing.balance.utils.ImapEvent;
|
||||||
import ru.spcex.clearing.balance.utils.MatcherFactory;
|
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
|
import ru.spcex.platform.enumeration.Task;
|
||||||
|
|
||||||
import javax.annotation.PostConstruct;
|
import javax.annotation.PostConstruct;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
|
|
@ -20,13 +18,13 @@ import java.util.Collection;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
|
|
||||||
import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFieldsComparator;
|
|
||||||
|
|
||||||
class Sdf08ServiceTest extends AbstractServiceTest {
|
class Sdf08ServiceTest extends AbstractServiceTest {
|
||||||
private static final MatcherFactory.Matcher<SDf08> SDF_08_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
// private static final MatcherFactory.Matcher<SDf08> SDF_08_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
||||||
|
private static final String TOPIC = Task.getAllBalance.topic();
|
||||||
|
private static final int PARTITION = 1;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
private Consumer<String, Object> mockConsumer;
|
private MockConsumer<String, Object> mockConsumer;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
Sdf08Service sdf08Service;
|
Sdf08Service sdf08Service;
|
||||||
|
|
@ -60,20 +58,17 @@ class Sdf08ServiceTest extends AbstractServiceTest {
|
||||||
|
|
||||||
//KAFKA
|
//KAFKA
|
||||||
mockConsumer.schedulePollTask(() -> {
|
mockConsumer.schedulePollTask(() -> {
|
||||||
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC_ACCOUNT_NEW, PARTITION)));
|
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC, PARTITION)));
|
||||||
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC_ACCOUNT_NEW, PARTITION, 0, "key", jsonBaseNewRequest));
|
mockConsumer.addRecord(new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonBaseNewRequest));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
||||||
HashMap<TopicPartition, Long> startOffsets = new HashMap<>();
|
HashMap<TopicPartition, Long> startOffsets = new HashMap<>();
|
||||||
TopicPartition tp = new TopicPartition(TOPIC_ACCOUNT_NEW, PARTITION);
|
TopicPartition tp = new TopicPartition(TOPIC, PARTITION);
|
||||||
startOffsets.put(tp, 0L);
|
startOffsets.put(tp, 0L);
|
||||||
mockConsumer.updateBeginningOffsets(startOffsets);
|
mockConsumer.updateBeginningOffsets(startOffsets);
|
||||||
|
|
||||||
IMap<Long, BankAccount> iMap = hazelcastServiceTest.getHazelcast().getMap(IMDGDistributedNames.Map_BankAccount);
|
|
||||||
|
|
||||||
//waiting for hazelcast map item updates
|
//waiting for hazelcast map item updates
|
||||||
ImapEvent imapEvent = new ImapEvent(iMap);
|
ImapEvent imapEvent = new ImapEvent(sdf08Map);
|
||||||
imapEvent.waitWhenHappened();
|
imapEvent.waitWhenHappened();
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue