This commit is contained in:
parent
05d015fd26
commit
6ad470b67f
3 changed files with 7 additions and 1 deletions
1
.gitignore
vendored
1
.gitignore
vendored
|
|
@ -2,3 +2,4 @@ target
|
||||||
/.idea/
|
/.idea/
|
||||||
*.iml
|
*.iml
|
||||||
logs
|
logs
|
||||||
|
tmp_notes.**
|
||||||
|
|
|
||||||
|
|
@ -11,6 +11,7 @@ public class MoneyMarketCreateAction implements IAction<MoneyMarketCreateRequest
|
||||||
@JsonFormat(pattern = "yyyy-MM-dd", timezone = "Europe/Moscow")
|
@JsonFormat(pattern = "yyyy-MM-dd", timezone = "Europe/Moscow")
|
||||||
@JsonProperty
|
@JsonProperty
|
||||||
public Instant startDate;
|
public Instant startDate;
|
||||||
|
@JsonFormat(pattern = "yyyy-MM-dd", timezone = "Europe/Moscow")
|
||||||
@JsonProperty
|
@JsonProperty
|
||||||
public Instant endDate;
|
public Instant endDate;
|
||||||
@JsonProperty
|
@JsonProperty
|
||||||
|
|
|
||||||
|
|
@ -5,8 +5,11 @@ 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.ConsumerRecords;
|
import org.apache.kafka.clients.consumer.ConsumerRecords;
|
||||||
import org.apache.kafka.common.errors.WakeupException;
|
import org.apache.kafka.common.errors.WakeupException;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
import ru.spcex.clearing.platform.messaging.logic.functional.BuilderConsumerStep;
|
import ru.spcex.clearing.platform.messaging.logic.functional.BuilderConsumerStep;
|
||||||
import ru.spcex.clearing.platform.messaging.logic.functional.ConsumerSpecificClass;
|
import ru.spcex.clearing.platform.messaging.logic.functional.ConsumerSpecificClass;
|
||||||
|
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||||
|
|
||||||
import java.time.Duration;
|
import java.time.Duration;
|
||||||
import java.time.temporal.ChronoUnit;
|
import java.time.temporal.ChronoUnit;
|
||||||
|
|
@ -17,6 +20,7 @@ import java.util.concurrent.Executors;
|
||||||
import java.util.concurrent.atomic.AtomicBoolean;
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
|
|
||||||
public class QueueConsumer implements AutoCloseable {
|
public class QueueConsumer implements AutoCloseable {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
private final AtomicBoolean closed = new AtomicBoolean(false);
|
private final AtomicBoolean closed = new AtomicBoolean(false);
|
||||||
private final Consumer<String, Object> consumer;
|
private final Consumer<String, Object> consumer;
|
||||||
private final ExecutorService executor;
|
private final ExecutorService executor;
|
||||||
|
|
@ -45,7 +49,7 @@ public class QueueConsumer implements AutoCloseable {
|
||||||
} catch (WakeupException e) {
|
} catch (WakeupException e) {
|
||||||
if (!closed.get()) throw e;
|
if (!closed.get()) throw e;
|
||||||
} catch (Throwable e) {
|
} catch (Throwable e) {
|
||||||
throw new RuntimeException(e);
|
log.error(ExceptionUtils.getStackTrace(e));
|
||||||
} finally {
|
} finally {
|
||||||
consumer.close();
|
consumer.close();
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue