【发布时间】:2021-10-31 22:34:07
【问题描述】:
所以,我正在使用 axon 和 Spring Boot 框架为低延迟交易引擎开发 PoC。单个流程流是否有可能实现低至 10 - 50 毫秒的延迟?该过程将包括验证、订单和风险管理。我已经对一个简单的应用程序进行了一些初步测试,以更新订单状态并执行它,我的延迟时间为 300 毫秒以上。这让我好奇我可以用 Axon 优化多少?
编辑:
延迟问题与 Axon 无关。使用InMemoryEventStorageEngine 和DisruptorCommandBus 设法将每个流程流降低到约5ms。
消息流是这样的。 NewOrderCommand(从客户端发布)-> OrderCreated(从聚合发布)-> ExecuteOrder(从 saga 发布)-> OrderExecutionRequested -> ConfirmOrderExecution(从 saga 发布)-> OrderExecuted(从聚合发布)
编辑 2: 最后切换到 Axon 服务器,但正如预期的那样,平均延迟上升到 ~150 毫秒。 Axon Server 是使用 Docker 安装的。如何使用 AxonServer 优化应用程序以实现亚毫秒级延迟?任何指针表示赞赏。
编辑 3:
@Steven,根据您的建议,我设法将延迟降低到平均 10 毫秒,这是一个好的开始!但是,是否有可能进一步降低它?因为我现在正在测试的只是在最终执行订单之前需要完成的一系列流程中的一个小流程,例如验证、风险管理和位置跟踪。所有这些都应在 5ms 或更短的时间内完成。最坏的情况是 10 毫秒(这些是更新的时间预算)。另外,请注意下面的配置,新读数基于由WeakReferenceCache 支持的InMemorySagaStore。非常感谢您的帮助!
订单聚合:
@Aggregate
internal class OrderAggregate {
@AggregateIdentifier(routingKey = "orderId")
private lateinit var clientOrderId: String
private var orderId: String = UUID.randomUUID().toString()
private lateinit var state: OrderState
private lateinit var createdAtSource: LocalTime
private val log by Logger()
constructor() {}
@CommandHandler
constructor(command: NewOrderCommand) {
log.info("received new order command")
val (orderId, created) = command
apply(
OrderCreatedEvent(
clientOrderId = orderId,
created = created
)
)
}
@CommandHandler
fun handle(command: ConfirmOrderExecutionCommand) {
apply(OrderExecutedEvent(orderId = command.orderId, accountId = accountId))
}
@CommandHandler
fun execute(command: ExecuteOrderCommand) {
log.info("execute order event received")
apply(
OrderExecutionRequestedEvent(
clientOrderId = clientOrderId
)
)
}
@EventSourcingHandler
fun on(event: OrderCreatedEvent) {
log.info("order created event received")
clientOrderId = event.clientOrderId
createdAtSource = event.created
setState(Confirmed)
}
@EventSourcingHandler
fun on(event: OrderExecutedEvent) {
val now = LocalTime.now()
log.info(
"elapse to execute: ${
createdAtSource.until(
now,
MILLIS
)
}ms. created at source: $createdAtSource, now: $now"
)
setState(Executed)
}
private fun setState(state: OrderState) {
this.state = state
}
}
OrderManagerSaga:
@Profile("rabbit-executor")
@Saga(sagaStore = "sagaStore")
class OrderManagerSaga {
@Autowired
private lateinit var commandGateway: CommandGateway
@Autowired
private lateinit var executor: RabbitMarketOrderExecutor
private val log by Logger()
@StartSaga
@SagaEventHandler(associationProperty = "clientOrderId")
fun on(event: OrderCreatedEvent) {
log.info("saga received order created event")
commandGateway.send<Any>(ExecuteOrderCommand(orderId = event.clientOrderId, accountId = event.accountId))
}
@SagaEventHandler(associationProperty = "clientOrderId")
fun on(event: OrderExecutionRequestedEvent) {
log.info("saga received order execution requested event")
try {
//execute order
commandGateway.send<Any>(ConfirmOrderExecutionCommand(orderId = event.clientOrderId))
} catch (e: Exception) {
log.error("failed to send order: $e")
commandGateway.send<Any>(
RejectOrderCommand(
orderId = event.clientOrderId
)
)
}
}
}
豆子:
@Bean
fun eventSerializer(mapper: ObjectMapper): JacksonSerializer{
return JacksonSerializer.Builder()
.objectMapper(mapper)
.build()
}
@Bean
fun commandBusCache(): Cache {
return WeakReferenceCache()
}
@Bean
fun sagaCache(): Cache {
return WeakReferenceCache()
}
@Bean
fun associationsCache(): Cache {
return WeakReferenceCache()
}
@Bean
fun sagaStore(sagaCache: Cache, associationsCache: Cache): CachingSagaStore<Any>{
val sagaStore = InMemorySagaStore()
return CachingSagaStore.Builder<Any>()
.delegateSagaStore(sagaStore)
.associationsCache(associationsCache)
.sagaCache(sagaCache)
.build()
}
@Bean
fun commandBus(
commandBusCache: Cache,
orderAggregateFactory: SpringPrototypeAggregateFactory<Order>,
eventStore: EventStore,
txManager: TransactionManager,
axonConfiguration: AxonConfiguration,
snapshotter: SpringAggregateSnapshotter
): DisruptorCommandBus {
val commandBus = DisruptorCommandBus.builder()
.waitStrategy(BusySpinWaitStrategy())
.executor(Executors.newFixedThreadPool(8))
.publisherThreadCount(1)
.invokerThreadCount(1)
.transactionManager(txManager)
.cache(commandBusCache)
.messageMonitor(axonConfiguration.messageMonitor(DisruptorCommandBus::class.java, "commandBus"))
.build()
commandBus.registerHandlerInterceptor(CorrelationDataInterceptor(axonConfiguration.correlationDataProviders()))
return commandBus
}
Application.yml:
axon:
server:
enabled: true
eventhandling:
processors:
name:
mode: tracking
source: eventBus
serializer:
general : jackson
events : jackson
messages : jackson
【问题讨论】:
-
我相信你可以,因为 Axon Framework 中的任何组件都可以配置。此外,您还可以将基础架构组件切换为更优化的部件。不过,在这个阶段,您的问题是广泛推荐您可以配置的 what。有关消息流的任何细节以及应用程序的建模方式都将在此处提供进一步的见解。所以,如果你能更新你的问题,那就太好了。
-
嗨@Steven,我已经更新了我的问题
-
感谢您扩展 David 的主题。
标签: axon