我使用带有 kafka 的 quarkus smallrye 反应式消息传递得到以下堆栈跟踪:
2020-07-24 01:38:31,662 ERROR [io.sma.rea.mes.kafka] (executor-thread-870) SRMSG18207: Unable to dispatch message to Kafka: io.smallrye.mutiny.subscription.BackPressureFailure: Could not emit tick 211 due to lack of requests
at io.smallrye.mutiny.operators.multi.builders.IntervalMulti$IntervalRunnable.run(IntervalMulti.java:83)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294)
at org.jboss.threads.ContextClassLoaderSavingRunnable.run(ContextClassLoaderSavingRunnable.java:35)
at org.jboss.threads.EnhancedQueueExecutor.safeRun(EnhancedQueueExecutor.java:2046)
at org.jboss.threads.EnhancedQueueExecutor$ThreadBody.doRunTask(EnhancedQueueExecutor.java:1578)
at org.jboss.threads.EnhancedQueueExecutor$ThreadBody.run(EnhancedQueueExecutor.java:1452)
at org.jboss.threads.DelegatingRunnable.run(DelegatingRunnable.java:29)
at org.jboss.threads.ThreadLocalResettingRunnable.run(ThreadLocalResettingRunnable.java:29)
at java.lang.Thread.run(Thread.java:748)
at org.jboss.threads.JBossThread.run(JBossThread.java:479)
所以我读过https://smallrye.io/smallrye-mutiny/#_how_do_i_control_the_back_pressure
根据文档,我添加了 BackPressure 控件。
前 :
@Outgoing( "eqs-crossing-xxx" )
public Multi< EQSAlert > eqsCrossingXXX_XXX(){
final String series = CrossingEnum.XXX_XXX.getSeries();
final String equipment = CrossingEnum.XXX_XXX.getEquipment();
final String vehicleRegex = vehicleRegexService.getRegexBySeries( series );
log.info( "Incoming request for {} - {}", series, equipment);
log.info( "Vehicle regex : {}", vehicleRegex );
return Multi
.createFrom()
.ticks()
.every(
Duration.ofSeconds( poolingInterval )
)
.concatMap(i -> {
final Multi<CrossingState> crossingStateBySeriesAndEquipment = CrossingState.getCrossingStateBySeriesAndEquipment(client, series, equipment);
return crossingStateBySeriesAndEquipment.flatMap(crossingState ->
crossingState.isActive() ?
EQSAlert.getEQSAlertBySeriesAndEquipment(
client,
series,
vehicleRegex,
equipment
)
:
Multi.createFrom().empty()
);
});
}
后 :
@Outgoing( "eqs-crossing-xxx" )
public Multi< EQSAlert > eqsCrossingXXX_XXX(){
final String series = CrossingEnum.XXX_XXX.getSeries();
final String equipment = CrossingEnum.XXX_XXX.getEquipment();
final String vehicleRegex = vehicleRegexService.getRegexBySeries( series );
log.info( "Incoming request for {} - {}", series, equipment);
log.info( "Vehicle regex : {}", vehicleRegex );
return Multi
.createFrom()
.ticks()
.every(
Duration.ofSeconds( poolingInterval )
)
.onOverflow()
.buffer(10)
.concatMap(i -> {
final Multi<CrossingState> crossingStateBySeriesAndEquipment = CrossingState.getCrossingStateBySeriesAndEquipment(client, series, equipment);
return crossingStateBySeriesAndEquipment.flatMap(crossingState ->
crossingState.isActive() ?
EQSAlert.getEQSAlertBySeriesAndEquipment(
client,
series,
vehicleRegex,
equipment
)
:
Multi.createFrom().empty()
);
});
}
现在一切都好了。
我的帖子的目的是了解我为什么需要这样做?
为什么缓冲区不能保留你?
如您所见,我每 5 秒(poolingInterval)调用一次简单的 sql 函数。该函数返回一些记录(通过池化从不超过 10 个)
所以流量非常低
请允许我用一些话来理解缓冲区管理。
谢谢