【问题标题】:Ultra low latency processes using Axon Framework使用 Axon 框架的超低延迟进程
【发布时间】:2021-10-31 22:34:07
【问题描述】:

所以,我正在使用 axon 和 Spring Boot 框架为低延迟交易引擎开发 PoC。单个流程流是否有可能实现低至 10 - 50 毫秒的延迟?该过程将包括验证、订单和风险管理。我已经对一个简单的应用程序进行了一些初步测试,以更新订单状态并执行它,我的延迟时间为 300 毫秒以上。这让我好奇我可以用 Axon 优化多少?

编辑:
延迟问题与 Axon 无关。使用InMemoryEventStorageEngineDisruptorCommandBus 设法将每个流程流降低到约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


【解决方案1】:

原始回复

您的设置描述很详尽,但我认为我仍然可以推荐一些选项。这涉及到框架中的许多位置,因此如果根据他们在 Axon 中的位置或目标对建议有任何不清楚的地方,请随时添加评论,以便我更新我的回复。

现在,让我们列出我的想法:

  • 为聚合设置快照如果加载需要很长时间。可使用AggregateLoadTimeSnapshotTriggerDefinition 进行配置。
  • 为您的聚合引入缓存。我会先尝试WeakReferenceCache。如果这还不够,那么值得研究 EhCache 和 JCache 适配器。或者,构建自己的。 Here's 顺便说一下聚合缓存部分。
  • 为您的传奇引入缓存。我会先尝试WeakReferenceCache。如果这还不够,那么值得研究 EhCache 和 JCache 适配器。或者,构建自己的。 Here's 顺便说一下,关于 Saga 缓存的部分。
  • 在此设置中您真的需要 Saga 吗?该过程看起来很简单,可以在常规事件处理组件中运行。如果是这种情况,在 Saga 流程中移动也可能会提高速度。
  • 您是否尝试过优化DisruptorCommandBus?尝试使用WaitStrategy、发布者线程数、调用者线程数和使用的Executor
  • 试试PooledStreamingEventProcessor(简称PSEP)而不是TrackingEventProcessor(简称TEP)。前者提供了更多的配置选项。顺便说一句,与 TEP 相比,默认值已经提供了更高的吞吐量。增加“批量大小”可以让您一次摄取更多的事件。您还可以更改 PSEP 用于事件检索工作(由 协调器 完成)和事件处理(worker 执行器负责此)的Executor
  • 您还可以在 Axon 服务器上配置一些可能会增加吞吐量的东西。试试event.events-per-segment-prefetchevent.read-buffer-sizecommand-thread。可能还有其他可用的选项,因此可能值得查看整个选项列表here
  • 虽然很难推断这是否会立即产生好处,但您可以为 Axon Server 运行更多内存/CPU。至少 2Gb 堆和 4 个核心。使用这些数字可能也会有所帮助。

可能还有更多要分享的内容,但这些是我的首要任务。希望这对您有所帮助,大卫!


第二反应

为了进一步推断我们可以在哪些方面获得更高的性能,我认为了解您的应用程序正在处理的哪个进程花费的时间最长是很重要的。如果我们可以改进的话,这将使我们能够推断出应该改进的地方。

您是否尝试过进行线程转储来推断哪个部分占用的时间最多?如果您可以将其分享为您的问题的更新,我们可以开始考虑以下步骤。

【讨论】:

  • 嗨,@Steven,感谢您的详尽回答,非常有帮助!从那以后,我设法测试了您的一些建议。我已经对我的发现进行了第三次编辑更新了这个问题。
  • 惊人的消息,你把它带到了稳定的 10 毫秒大卫!我更新了我的帖子,为您添加了另一个问题,这可能会导致您的更新。
  • 嗨@Steven,终于来捕获线程转储了。试图在 fastthread.io 上分析它,但不太确定要寻找什么。感谢你能指出我正确的方向
  • 据我所知,我会检查哪些类/方法占用了线程的大部分时间。排除任何 JDK 默认值并向下移动,直到找到 (1) 您的代码或 (2) Axon 代码。如果是 Axon 代码,我想我可以建议是否有更快的方法。
  • 不,因为TokenStore 旨在用作应用程序实例之间的通信方式。因此,缓存实现需要成为分布式解决方案,我们现阶段不打算实现。
猜你喜欢
  • 2018-03-01
  • 2022-09-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-01-23
  • 2016-01-24
  • 2020-10-11
  • 1970-01-01
相关资源
最近更新 更多