Python实现LMAX架构Disruptor模式 业务逻辑层数据读取实现咨询
LMAX架构Disruptor模式Python实现相关问题
我正尝试在Lmax架构中实现Disruptor模式。众所周知,LMAX架构中通过Ring Buffer构建队列来处理数据,其结构如下所示:
我已使用Python实现了该结构,代码如下:
import multiprocessing class CircularBuffer(object): def __init__(self, max_size=10): """Initialize the CircularBuffer with a max_size if set, otherwise max_size will elementsdefault to 10""" self.buffer = [None] * max_size self.blconsumer = 0 self.receiver = 0 self.journalerPointer=0 self.replicatorPointer=0 self.unmarshallerPointer=0 self.max_size = max_size def __str__(self): """Return a formatted string representation of this CircularBuffer.""" items = ['{!r}'.format(item) for item in self.buffer] return '[' + ', '.join(items) + ']' def size(self): """Return the size of the CircularBuffer Runtime: O(1) Space: O(1)""" if self.receiver >= self.blconsumer: return self.receiver - self.blconsumer return self.max_size - self.head - self.receiver def is_empty(self): """Return True if the head of the CircularBuffer is equal to the tail, otherwise return False Runtime: O(1) Space: O(1)""" return self.receiver == self.blconsumer def is_replicator_after_receiver(self): """Return True if the head of the CircularBuffer is equal to the tail, otherwise return False Runtime: O(1) Space: O(1)""" return self.receiver == (self.replicatorPointer-1) % self.max_size def is_journaler_after_receiver(self): """Return True if the head of the CircularBuffer is equal to the tail, otherwise return False Runtime: O(1) Space: O(1)""" return self.receiver == (self.journalerPointer-1) % self.max_size def is_unmarshaller_after_receiver(self): """Return True if the head of the CircularBuffer is equal to the tail, otherwise return False Runtime: O(1) Space: O(1)""" return self.receiver == (self.unmarshallerPointer-1) % self.max_size def is_BusinessLogicConsumer_after_unmarshaller(self): """Return True if the head of the CircularBuffer is equal to the tail, otherwise return False Runtime: O(1) Space: O(1)""" return self.unmarshallerPointer == (self.blconsumer-1) % self.max_size def is_full(self): """Return True if the tail of the CircularBuffer is one before the head, otherwise return False Runtime: O(1) Space: O(1)""" return self.receiver == (self.blconsumer-1) % self.max_size def receive(self, item): """Insert an item at the back of the CircularBuffer Runtime: O(1) Space: O(1)""" if self.is_full()==False : self.buffer[self.receiver] = item self.receiver = (self.receiver + 1) % self.max_size def front(self): """Return the item at the front of the CircularBuffer Runtime: O(1) Space: O(1)""" return self.buffer[self.blconsumer] def consume(self): """Return the item at the front of the Circular Buffer and remove it Runtime: O(1) Space: O(1)""" # if self.is_empty(): # raise IndexError("CircularBuffer is empty, unable to dequeue") # if self.is_BusinessLogicConsumer_after_unmarshaller()==True : # raise IndexError("BusinessLogicConsumer can't be after receiver") if self.is_BusinessLogicConsumer_after_unmarshaller()==False and self.is_empty()==False: item = self.buffer[self.blconsumer] self.buffer[self.blconsumer] = None self.blconsumer = (self.blconsumer + 1) % self.max_size return item def replicator(self): # if self.is_empty(): # raise IndexError("CircularBuffer is empty, unable to dequeue") # if self.is_replicator_after_receiver()==True : # raise IndexError("replicator can't be after receiver") if self.is_replicator_after_receiver()==False and self.is_empty()==False: item = self.buffer[self.replicatorPointer] self.replicatorPointer = (self.replicatorPointer + 1) % self.max_size return item def journaler(self): # if self.is_empty(): # raise IndexError("CircularBuffer is empty, unable to dequeue") # if self.is_journaler_after_receiver()==True : # raise IndexError("journaler can't be after receiver") if self.is_journaler_after_receiver()==False and self.is_empty()==False: item = self.buffer[self.journalerPointer] self.journalerPointer = (self.journalerPointer + 1) % self.max_size return item def unmarshaller(self): # if self.is_empty(): # raise IndexError("CircularBuffer is empty, unable to dequeue") # if self.is_unmarshaller_after_receiver()==True : # raise IndexError("unmarshaller can't be after receiver") if self.is_unmarshaller_after_receiver()==False and self.is_empty()==False: item = self.buffer[self.journalerPointer] self.unmarshallerPointer = (self.unmarshallerPointer + 1) % self.max_size return item
如架构图所示,LMAX的业务逻辑模块会从Ring Buffer拉取数据到CPU进行高速处理,但目前我没有找到实现业务逻辑层的相关文档,想咨询在Python中实现LMAX架构时,如何将数据从Ring Buffer读取到CPU寄存器中完成业务逻辑处理?
回答
首先明确一个基础认知:Python是上层解释型语言,没有提供直接操控CPU寄存器的底层接口,寄存器的数据加载调度由CPython解释器、操作系统共同完成。LMAX原生Java实现的高性能也并非手动操作寄存器实现,而是通过无锁设计、缓存行对齐、避免GC停顿等优化,尽可能减少CPU等待开销,让数据可以长期命中CPU高速缓存,接近寄存器级别的访问效率。
Python环境下可以按照以下方案实现符合LMAX设计理念的高性能业务逻辑层:
- 先修复现有Ring Buffer的已知问题
- size方法中引用的
self.head未定义,应该替换为self.blconsumer - unmarshaller方法中读取的是
self.journalerPointer位置的数据,不符合多消费者独立指针的设计,应该改为读取self.unmarshallerPointer对应位置的数据
- size方法中引用的
- 业务消费线程做核心绑定
遵循LMAX单线程业务处理的核心设计,将业务逻辑消费者单独绑定到固定CPU核心运行,减少进程上下文切换开销,Python中可以通过psutil.Process().cpu_affinity([core_id])实现核心绑定。 - 批量拉取数据处理
不要每次调用consume只拉取一条数据,每次消费时先获取当前Ring Buffer中所有可用的可消费数据,批量执行业务逻辑,减少方法调用开销,同时提升CPU缓存命中率。 - 优化业务逻辑本身
业务处理逻辑中不要加入任何磁盘IO、网络请求等阻塞操作,也不要频繁创建临时大对象,减少GC停顿带来的性能损耗。 - 优先使用成熟的第三方Disruptor实现
开源社区已经有成熟的Python版Disruptor实现,内置了缓存行填充、无锁序列等优化,性能远高于手写的环形缓冲区,不需要重复造轮子。
简易业务层实现示例
import psutil import os # 业务逻辑处理类 class BusinessProcessor: def process(self, data): # 此处写入你的业务处理逻辑,确保是纯内存计算,无阻塞操作 return data * 2 def business_consumer_loop(buffer: CircularBuffer, processor: BusinessProcessor): # 绑定到第3个CPU核心,可根据实际情况调整 p = psutil.Process(os.getpid()) p.cpu_affinity([2]) while True: # 批量拉取可用数据 batch_data = [] while not buffer.is_BusinessLogicConsumer_after_unmarshaller() and not buffer.is_empty(): item = buffer.consume() if item: batch_data.append(item) # 批量处理数据 for data in batch_data: processor.process(data)
内容的提问来源于stack exchange,提问作者Ehsan Akbar
相关产品推荐
相关产品推荐

