LMAX Disruptor通常使用以下方法实现:

在此示例中,Replicator负责将输入事件\命令复制到从属节点.复制一组节点需要我们应用一致性算法,以防我们希望系统在出现网络故障,主故障和从站故障时可用.
我正在考虑将RAFT一致性算法应用于此问题.一个观察结果是:"RAFT要求在复制期间将输入事件\命令存储到磁盘(持久存储)"(参考此链接)
这种观察实质上意味着我们无法执行内存中复制.因此,似乎我们可能必须结合复制器和记者的功能才能成功地将RAFT算法应用于LMAX.
有两种方法可以做到这一点:
选项1:使用复制日志作为输入事件队列

我认为这个选项的缺点与我们做一个额外的数据复制步骤(接收器到事件队列而不是环形缓冲区)这一事实有关.
选项2:使用Replicator将输入事件\命令推送到从属的输入日志文件

我想知道是否有其他解决方案来设计Replicator?人们用于复制器的不同设计选择有哪些?特别是任何可以支持内存复制的设计?
我最近一直在学习LMAX Disruptor并且正在做一些实验.令我困惑的一件事是处理程序方法的endOfBatch参数.请考虑以下代码.首先,我调用的虚拟消息和消费者类,以及:onEventEventHandlerTest1Test1Worker
public class Test1 {
}
public class Test1Worker implements EventHandler<Test1>{
public void onEvent(Test1 event, long sequence, boolean endOfBatch) {
try{
Thread.sleep(500);
}
catch(Exception e){
e.printStackTrace();
}
System.out.println("Received message with sequence " + sequence + ". "
+ "EndOfBatch = " + endOfBatch);
}
}
Run Code Online (Sandbox Code Playgroud)
请注意,我已经延迟了500毫秒,以替代一些现实世界的工作.我也在控制台中打印了序列号
然后我的驱动程序类(作为生产者)调用DisruptorTest:
public class DisruptorTest {
private static Disruptor<Test1> bus1;
private static ExecutorService test1Workers;
public static void main(String[] args){
test1Workers = Executors.newFixedThreadPool(1);
bus1 = new Disruptor<Test1>(new …Run Code Online (Sandbox Code Playgroud) java multithreading producer-consumer disruptor-pattern lmax
我从数据馈送 (lmax) 中获取一次分钟数据 (OHLC),我想将其重新采样为五分钟数据,然后再将其重新采样为十分钟、十五分钟和三十分钟。
我正在使用以下逻辑:
开盘=每 5 根蜡烛的第一个值(第一根蜡烛打开)
高 = 5 根蜡烛的最高值
低 = 5 根蜡烛的最低值
Close=最后一根蜡烛的收盘价(第 5 根蜡烛)
这是重新采样数据的正确方法吗?我觉得这是合乎逻辑的,但出于某种原因,他们网站上的数据源和我的代码之间存在明显的差异。我觉得是这样,因为我重新采样错了;如果我使用 Python,我可以参考 Pandas,但不能使用 C#(据我所知)。
这是重新采样数据的函数:
private static List<List<double>> normalize_candles(List<List<double>> indexed_data, int n)
{
double open = 0;
double high = 0;
double low = 5;
double close = 0;
int trunc = 0;
if (indexed_data.Count() % n != 0)
{
trunc = indexed_data.Count() % n;
for (int i = 0; i < trunc; i++)
{
indexed_data.RemoveAt(indexed_data.Count() - 1);
}
}
{ …Run Code Online (Sandbox Code Playgroud) 如何监控LMAX Disruptor?假设我有 3 个环形缓冲区,并希望提供一个 ui 来为我提供环形缓冲区的信息。
遵循Disruptor入门指南,我建立了一个由单个生产者和单个消费者组成的最小破坏者。
制片人
import com.lmax.disruptor.RingBuffer;
public class LongEventProducer
{
private final RingBuffer<LongEvent> ringBuffer;
public LongEventProducer(RingBuffer<LongEvent> ringBuffer)
{
this.ringBuffer = ringBuffer;
}
public void onData()
{
long sequence = ringBuffer.next();
try
{
LongEvent event = ringBuffer.get(sequence);
}
finally
{
ringBuffer.publish(sequence);
}
}
}
Run Code Online (Sandbox Code Playgroud)
消费者(注意消费者什么也不做onEvent)
import com.lmax.disruptor.EventHandler;
public class LongEventHandler implements EventHandler<LongEvent>
{
public void onEvent(LongEvent event, long sequence, boolean endOfBatch)
{}
}
Run Code Online (Sandbox Code Playgroud)
我的目标是对大型环缓冲区进行一次性能测试,而不是多次遍历较小的环。在每种情况下,总操作数(bufferSizeX rotations)都是相同的。我发现随着环形缓冲区变小,操作/秒速率急剧下降。
RingBuffer Size | Revolutions | Total Ops | Mops/sec
1048576 …Run Code Online (Sandbox Code Playgroud) 请考虑Martin Fowler的 LMAX Architecture 描述中的以下场景:
我将使用一个简单的非LMAX示例来说明.想象一下,您正在通过信用卡订购果冻豆.<...>
在LMAX架构中,您可以将此操作拆分为两个.第一个操作将捕获订单信息并通过向信用卡公司输出事件(请求的信用卡验证)来完成.然后,业务逻辑处理器将继续为其他客户处理事件,直到它在其输入事件流中收到信用卡验证事件.在处理该事件时,它将执行该订单的确认任务.
因此,订单保留在内存中,直到收到付款处理结果.
现在让我们假设代替信用卡处理步骤,我们需要花费更多时间的步骤,例如:我们需要执行库存检查,有人必须在物理上验证我们是否已经订购了特定的果冻豆味道.这可能需要一个小时.
如果是这种情况,不会导致内存中保存的数据增长,因为很多订单可能会等待库存状态更新事件?
可能在这种情况下,我们需要从内存中删除订单并将其作为输出事件的一部分包含在内,外部系统(库存)负责生成包含订单详细信息的另一个输入事件.
我用这种方法看到的问题是,我们不能将库存作为业务逻辑处理器的一部分.
关于我们如何解决这个问题的想法?
在LMAX Disruptor模式中,复制器用于将输入事件从主节点复制到从节点.所以设置可能如下所示:

主节点的复制器将事件写入数据库(尽管我们可以考虑比写入数据库更好的机制 - 但它对问题语句并不重要).从节点的接收器从DB读取并将事件放入从节点的环形缓冲器.
从节点的输出事件被忽略.
现在,主节点的业务逻辑处理器可能比从节点的业务逻辑处理器慢.例如,主节点的BL可以在时隙102处,其中从节点可以在106处.(这可以发生,因为复制器在业务逻辑处理器之前从环形缓冲器读取事件).
在这种情况下,如果主节点发生故障并且从节点现在成为主节点,则外部系统可能会遗漏一些关键事件.这可能发生,因为节点2在充当从节点时会忽略其输出.
Martin Fowler确实说复制器的工作是保持节点同步:"之前我提到LMAX在集群中运行其系统的多个副本以支持快速故障转移.复制器使这些节点保持同步"
但我不确定它如何使Business Logic Processor保持同步?有任何想法吗?
我计划在我的破坏者中有许多并行消费者。
我需要每个消费者只消费为他们准备的消息。
例如,我有 A、B、C 类型的消息,我有类似的缓冲区
#1 - type A, #2 - type B, #3 - type C, #4 - type A, #5 - type C, #6 - type C, (and so on)
Run Code Online (Sandbox Code Playgroud)
我有每种类型的消费者。我怎样才能让 A 的消费者接受消息 1 和 4,对于类型 B - 消息 2,C - 消息 3、5、6?
重要提示:我希望处理是独立的。消费者不应该被链接起来,每个人都独立地在缓冲区中移动。如果 A 的消费者比 C 的消费者慢,则“类型 C”消费者对 #6 的处理可能会早于类型 A 的 #1 参与。
我很欣赏如何使用 LMAX 干扰器配置来做到这一点的解释。
我理解了的LMAX干扰器是,它是一个完整的吓人,快速,可怕的并发Java代码JAR,允许吞吐量每秒20万条(如果正确使用).
我们目前有一个ActiveMQ实例,它在我们需要的整个过程中很慢,大约每秒400条消息.我想知道我们是否会受益于重构我们的代码以使用LMAX,但是有以下问题:
并且,如果我完全偏离所有这些,并且似乎完全误解了LMAX干扰器的使用,那么有人可以提供何时使用它的具体示例?提前致谢!
我有以下关于破坏者的问题:
c1
P1 - c2 - c4 - c5
c3
其中c1到c3可以在p1之后并行工作,C4和C5在他们之后工作。
所以通常我会有这样的东西(P1和C1-C5是可运行的/可调用的)
p1.start();
p1.join();
c1.start();
c2.start();
c3.start();
c1.join();
c2.join();
c3.join();
c4.start();
c4.join();
c5.start();
c5.join();
Run Code Online (Sandbox Code Playgroud)
但是在 Disruptor 的情况下,我的事件处理程序都没有实现 Runnable 或 Callable,那么 Disruptor 框架最终是如何并行运行它们的呢?
采取以下情景:
我的消费者 C2 需要对事件进行一些注释的 web 服务调用,在 SEDA 中,我可以为这样的 10 个 C2 请求启动 10 个线程[用于将消息从队列中拉出 + 进行 Webservice 调用并更新下一个 SEDA 队列],这将确保我不会为 10 个请求中的每一个依次等待 Web 服务响应,在这种情况下,我的事件处理器 C2(如果)作为单个实例将依次等待 10 个 C2 请求。
lmax ×10
java ×5
concurrency ×2
architecture ×1
c# ×1
failover ×1
middleware ×1
performance ×1
raft ×1
replication ×1
resampling ×1