add client hz connection for imdg-hist

This commit is contained in:
etreschenkov 2025-08-08 16:34:54 +03:00
parent 63f94d6921
commit 60d9801a5f
9 changed files with 336 additions and 2 deletions

View file

@ -0,0 +1,28 @@
package ru.spcex.clearing.historyimdg.component;
import java.util.HashMap;
import java.util.Map;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.misc.SessionHistory;
import ru.spcex.clearing.historyimdg.component.builder.ISearchProxyBuilder;
import ru.spcex.clearing.historyimdg.component.builder.SearchProxySessionHistoryBuilder;
import ru.spcex.clearing.historyimdg.index.SearchProxy;
@Component
public class SearchProxyBuilder {
private final Map<Class<?>, ISearchProxyBuilder<? extends SearchProxy>> builderByClass;
public SearchProxyBuilder() {
builderByClass = new HashMap<>();
builderByClass.put(SessionHistory.class, new SearchProxySessionHistoryBuilder());
}
public SearchProxy build(Object obj, Class<?> clazz) {
ISearchProxyBuilder<? extends SearchProxy> specificBuilder = builderByClass.get(clazz);
if (specificBuilder == null) {
throw new IllegalStateException("cannot construct SearchProxy for class " + clazz.getSimpleName());
}
return specificBuilder.build(obj);
}
}

View file

@ -0,0 +1,76 @@
package ru.spcex.clearing.historyimdg.component;
import com.hazelcast.core.EntryEvent;
import com.hazelcast.core.IMap;
import com.hazelcast.map.listener.EntryAddedListener;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
import org.apache.commons.lang3.exception.ExceptionUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.spcex.clearing.historyimdg.index.SearchProxy;
import ru.spcex.platform.classes.base.interfaces.WithId;
/**
* слушатель, который копит новые записи в указанной мапе<br>
* при вызове startProcessingNewEntries() начинает запись новых SearchProxy в мапу поиска
*/
public class TodayQueueListener<T extends WithId> implements EntryAddedListener<Long, T> {
private final Logger log = LoggerFactory.getLogger(getClass());
private IMap<Long, SearchProxy> targetMap;
private final Class<T> clazz;
private final String classNameForLog;
private final BlockingQueue<T> addedEntry;
private final ExecutorService executorService;//private boolean shutdownFlag = false;
private final SearchProxyBuilder searchProxyBuilder;
public TodayQueueListener(IMap<Long, SearchProxy> targetMap, Class<T> clazz, SearchProxyBuilder searchProxyBuilder) {
this.targetMap = targetMap;
this.clazz = clazz;
this.addedEntry = new LinkedBlockingQueue<>();
this.executorService = Executors.newSingleThreadExecutor();
this.classNameForLog = clazz.getSimpleName();
this.searchProxyBuilder = searchProxyBuilder;
}
@Override
public void entryAdded(EntryEvent<Long, T> event) {
T value = event.getValue();
if (value == null) {
log.warn("new entry {} received as null", clazz.getSimpleName());
return;
}
addedEntry.add(value);
}
public void startProcessingNewEntries() {
log.info("Start processing new entries for class {}", clazz.getSimpleName());
executorService.execute(() -> {
T newEntry;
while (true) { //!shutdownFlag
try {
newEntry = addedEntry.take();
Long key = newEntry.getId();
// if (targetMap.containsKey(key)) {
// log.trace("new entry for {} with id={} - already was downloaded", classNameForLog, key);
// continue;
// }
SearchProxy searchProxy = searchProxyBuilder.build(newEntry, clazz);
log.trace("new entry for {} with id={} constructed proxy: \n{}", classNameForLog, key, searchProxy);
targetMap.set(key, searchProxy);
} catch (InterruptedException e) { //значит мы хотим остановить этот обработчик
log.info("stopped processing new entries for class {}", classNameForLog);
return; //shutdownFlag = true;
} catch (Exception e) {
log.error(ExceptionUtils.getStackTrace(e));
}
}
});
}
public void shutdown() {
executorService.shutdownNow(); //shutdownFlag = true;
}
}

View file

@ -0,0 +1,8 @@
package ru.spcex.clearing.historyimdg.component.builder;
import ru.spcex.clearing.historyimdg.index.SearchProxy;
public interface ISearchProxyBuilder<T extends SearchProxy> {
T build(Object obj);
}

View file

@ -0,0 +1,16 @@
package ru.spcex.clearing.historyimdg.component.builder;
import ru.clearing.classes.statics.data.misc.SessionHistory;
import ru.spcex.clearing.historyimdg.index.SearchProxySessionHistory;
public class SearchProxySessionHistoryBuilder implements ISearchProxyBuilder<SearchProxySessionHistory> {
@Override
public SearchProxySessionHistory build(Object obj) {
SessionHistory value = (SessionHistory) obj;
SearchProxySessionHistory searchProxy = new SearchProxySessionHistory();
searchProxy.setId(value.getId());
searchProxy.setSessionId(value.getObject().getId());
return searchProxy;
}
}

View file

@ -0,0 +1,54 @@
package ru.spcex.clearing.historyimdg.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import ru.spcex.clearing.historyimdg.config.element.ImdgSettings;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@Configuration
public class HistoryHazelcastClientConfig {
Logger log = LoggerFactory.getLogger(getClass());
@Bean(name = "taskExecutorHazelcastClientInitializer")
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
return createThreadPoolTaskExecutor(1, true);
}
@Bean(name = "taskExecutorIdGeneratorAwaiter")
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
return createThreadPoolTaskExecutor(1, false);
}
@Autowired
@Bean
public ImdgProvider imdgProvider(
@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
@Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
ImdgSettings imdgSettings
) {
if (imdgSettings.getHazelcastClient() == null || imdgSettings.getHazelcastClient().getClusterMembers() == null) {
log.warn("Property \"backend-api.hazelcast.cluster-members\" not set!");
throw new IllegalArgumentException("Property \"backend-api.hazelcast.cluster-members\" not set");
}
ImdgProvider imdg = new HazelcastService(taskExecutorHazelcastClientInitializer,
taskExecutorIdGeneratorAwaiter,
imdgSettings.getHazelcastClient());
return imdg;
}
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
if (maxPoolSz > 2) {
pool.setKeepAliveSeconds(60);
pool.setAllowCoreThreadTimeOut(true);
}
pool.setCorePoolSize(maxPoolSz);
pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
return pool;
}
}

View file

@ -3,15 +3,25 @@ package ru.spcex.clearing.historyimdg.config.element;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
@PropertySource("file:${spring.config.location}/application.properties")
@ConfigurationProperties("imdg.hist")
public class ImdgSettings {
private HazelcastClientParams hazelcastClient;
private HazelcastServerSettings hazelcast;
private DatabaseSettings database;
private ControllerSettings debugServer;
public HazelcastClientParams getHazelcastClient() {
return hazelcastClient;
}
public void setHazelcastClient(HazelcastClientParams hazelcastClient) {
this.hazelcastClient = hazelcastClient;
}
public HazelcastServerSettings getHazelcast() {
return hazelcast;
}

View file

@ -42,6 +42,10 @@ public class HazelcastLifecycleSupport implements InitializingBean, DisposableBe
HazelcastHelper.imdgSystem_setStorageState(true, hazelcastServerInstance);
}
public HazelcastInstance getHazelcastServerInstance() {
return hazelcastServerInstance;
}
@Override
public void destroy() throws Exception {

View file

@ -0,0 +1,134 @@
package ru.spcex.clearing.historyimdg.services;
import com.hazelcast.aggregation.Aggregators;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.IMap;
import com.hazelcast.projection.Projections;
import com.hazelcast.query.Predicates;
import java.time.Instant;
import java.time.temporal.ChronoUnit;
import java.util.Collection;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.historyimdg.component.SearchProxyBuilder;
import ru.spcex.clearing.historyimdg.component.TodayQueueListener;
import ru.spcex.clearing.historyimdg.index.SearchProxy;
import ru.spcex.platform.classes.base.interfaces.WithId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.service.IHazelcastClusterStatus;
@Service
public class TodayHistoryService implements InitializingBean, DisposableBean, IHazelcastClusterStatus {
private final Logger log = LoggerFactory.getLogger(getClass());
private final ImdgProvider imdgProvider;
private final HazelcastLifecycleSupport hzLifecycleSupport;
private final SearchProxyBuilder searchProxyBuilder;
private HazelcastInstance availableStorageHzInstance;
//---- listeners ----
@Autowired
public TodayHistoryService(ImdgProvider imdgProvider, HazelcastLifecycleSupport hzLifecycleSupport, SearchProxyBuilder searchProxyBuilder) {
this.imdgProvider = imdgProvider;
this.hzLifecycleSupport = hzLifecycleSupport;
this.searchProxyBuilder = searchProxyBuilder;
}
@Override
public void afterPropertiesSet() {
// imdgProvider.statusSubscribe(this);
}
/**
* хазелкастовый серверный инстанс Search
*/
private HazelcastInstance searchHzServer() {
return hzLifecycleSupport.getHazelcastServerInstance();
}
/**
* регистрирует нового слушателя который будет сохранять индекс записей из стораджа в мапу для поиска
*
* @param storageMapName мапа с данными из стораджа
* @param searchMapName мапа с поисковыми индексами
*/
private <T extends WithId>
TodayQueueListener<T> registerListener(String storageMapName, String searchMapName, Class<T> clazz) {
IMap<Long, T> storageMap = availableStorageHzInstance.getMap(storageMapName);
TodayQueueListener<T> listener = new TodayQueueListener<>(searchHzServer().getMap(searchMapName), clazz, searchProxyBuilder);
storageMap.addEntryListener(listener, true);
return listener;
}
/**
* начинает поиск записей в storageMapName с id больше максимального id из searchMapName
* все записи трансформирует в SearchProxy и сохраняет в мапу поиска
*
* @param storageMapName мапа с данными из стораджа
* @param searchMapName мапа с поисковыми индексами
*/
public <T extends WithId>
void loadSnapshot(String storageMapName, String searchMapName, Class<T> clazz) {
Instant startTime = Instant.now();
IMap<Long, SearchProxy> targetMap = searchHzServer().getMap(searchMapName);
IMap<Long, T> historyMap = availableStorageHzInstance.getMap(storageMapName);
Long maxSearchId = targetMap.aggregate(Aggregators.longMax("id"));
if (maxSearchId == null) maxSearchId = 0L;
Collection<Long> ids = historyMap.project(Projections.singleAttribute("id")
, Predicates.greaterThan("id", maxSearchId));
for (Long id : ids) {
T t = historyMap.get(id);
targetMap.set(id, searchProxyBuilder.build(t, clazz));
}
log.info("downloaded {} today records from {} into search index map in {} ms with predicate id > {}",
ids.size(), storageMapName, ChronoUnit.MILLIS.between(startTime, Instant.now()), maxSearchId);
}
/**
* для каждой мапы
* <ol>
* <li>сначала регистрирует сохраняющего новые значения слушателя</li>
* <li>выгружает в индексовую мапу все значения с ID > текущего максимального</li>
* <li>запускает слушателей на запись</li>
* </ol>
* с таким подходом слушатель может по второму разу обработать записи которые появились после его регистрации,
* во время шага 2, эти записи он пропустит
*/
@Override
public void getAvailable(HazelcastInstance storageHzInstanceInited) {
this.availableStorageHzInstance = storageHzInstanceInited;
//-----
// repoMMOrderListener = registerListener(Map_RepoMMOrder, Map_RepoMMOrderSearch, RepoMMOrder.class);
// loadSnapshot(Map_RepoMMOrder, Map_RepoMMOrderSearch, RepoMMOrder.class);
// repoMMOrderListener.startProcessingNewEntries();
//-----
this.availableStorageHzInstance = null;
}
@Override
public void getUnavailable(HazelcastInstance hazelcastNotInited) {
//останавливаем всех слушателей
stopAllListeners();
}
@Override
public void destroy() throws Exception {
stopAllListeners();
}
protected void stopAllListeners() {
// repoMMOrderListener.shutdown();
}
public interface ConditionForEntity<T extends WithId> {
boolean verify(T item);
}
}

View file

@ -1,10 +1,14 @@
imdg.hist.hazelcast-client.cluster-members=10.200.200.183:5701
imdg.hist.hazelcast-client.login=dev
imdg.hist.hazelcast-client.password=dev-pass
imdg.hist.hazelcast.listenPort=5702
imdg.hist.hazelcast.login=dev-hist
imdg.hist.hazelcast.password=dev-pass-hist
imdg.hist.hazelcast.cluster-members[0]=127.0.0.1
imdg.hist.database.login=cls_dev
imdg.hist.database.login=clearing
imdg.hist.database.password=Aa111111
imdg.hist.database.url=jdbc:postgresql://10.200.200.133:5432/clearing?currentSchema=clearing_dev
imdg.hist.database.url=jdbc:postgresql://10.200.200.133:5432/clearing_uat
#imdg.database.url=jdbc:postgresql://10.200.200.133:5432/postgres?currentSchema=clearing_tester
#debug tester mode: