You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Thespian框架Actor间低延迟消息传输性能低下问题排查

Performance Lag Between Thespian Actors Sending/Receiving Messages

I'm using the Thespian framework to build an app that consumes low-latency data via Socket. Initially, I created two Actors: one as a consumer (Socket client connection) and another as a processor. During testing, I noticed a significant delay between when the consumer (Actor1) sends messages and when the processor (Actor2) receives them.

To test this issue, I used the following code:

from thespian.actors import ActorSystem, Actor
import time

class Actor1(Actor):
    def receiveMessage(self, handler, sender):
        # benchmark
        data = {'tic': time.perf_counter(), 'lim': 10000}
        print('Elapsed Time to process {} messages'.format(data['lim']))
        for i in range(data['lim']):
            self.send(handler, data)
        # perf
        toc = time.perf_counter() - data['tic']
        print('Actor 1: {} sec'.format(round(toc, 3)))
        # self.actorSystemShutdown()

class Actor2(Actor):
    def __init__(self):
        self.msg_count = 0
    def receiveMessage(self, data, sender):
        self.msg_count += 1
        if self.msg_count == data['lim']:
            toc = time.perf_counter() - data['tic']
            print('Actor 2: {} sec'.format(round(toc, 3)))

def main():
    asys = ActorSystem('multiprocTCPBase')
    consumer = asys.createActor(Actor1)
    handler = asys.createActor(Actor2)
    asys.tell(consumer, handler)

if __name__ == '__main__':
    main()

My test results are as follows:

Elapsed Time to process 10 messages
Actor 1: 0.002 sec
Actor 2: 0.008 sec

Elapsed Time to process 100 messages
Actor 1: 0.019 sec
Actor 2: 0.099 sec

Elapsed Time to process 1000 messages
Actor 1: 0.131 sec
Actor 2: 0.769 sec

Elapsed Time to process 10000 messages
Actor 1: 1.219 sec
Actor 2: 7.608 sec

Elapsed Time to process 100000 messages
Actor 1: 22.012 sec
Actor 2: 91.419 sec

I have three questions:

  1. Is there a problem or oversight in my code?
  2. Are there faster ways to send messages between Actors?
  3. Are there other performance benchmarking methods mentioned in the documentation to analyze this performance issue?

Great question—let’s break this down step by step.

1. Issues/Oversights in Your Current Code

First, let’s unpack what’s driving that growing lag:

  • Per-Message Serialization Overhead: You’re sending the exact same data dictionary thousands of times individually. Each message needs to be serialized for TCP transport and deserialized by Actor2, which adds up exponentially with message volume.
  • TCP Transport Overhead: multiprocTCPBase uses TCP sockets, which have inherent overhead for small, frequent messages (like packet framing, handshake overhead, and acknowledgment delays). Flooding the transport with tiny messages amplifies this cost.
  • Message Queue Backlog: Actors process messages one at a time. When Actor1 dumps 10k+ messages into Actor2’s queue all at once, the backlog builds up—this is why the gap between Actor1’s finish time and Actor2’s completion time grows linearly with message count.

2. Faster Actor-to-Actor Message Sending

Here are concrete optimizations to slash latency:

  • Batch Messages: Instead of sending thousands of individual messages, send a single batch (or chunks of 100-1000 messages at a time). This cuts serialization and transport overhead drastically. Example adjustment:
    # In Actor1:
    data = {'tic': time.perf_counter(), 'lim': 10000}
    batch = [data for _ in range(data['lim'])]
    self.send(handler, batch)
    
    # In Actor2:
    def receiveMessage(self, batch, sender):
        self.msg_count += len(batch)
        if self.msg_count == batch[0]['lim']:
            toc = time.perf_counter() - batch[0]['tic']
            print('Actor 2: {} sec'.format(round(toc, 3)))
    
  • Switch to a Lower-Overhead Transport: If your Actors run on the same machine, use multiprocQueueBase instead of multiprocTCPBase. It uses in-memory queues (no TCP serialization/network overhead) which is way faster for single-host inter-Actor communication. For cross-machine needs, consider multiprocUDPBase (just note UDP’s reliability tradeoffs for low-latency use cases).
  • Trim Unnecessary Message Data: You’re sending the entire data dict (including lim and tic) every time. Instead, send only the critical payload, and pass metadata like lim once upfront in a separate message.

3. Thespian Performance Benchmarking Methods

Thespian has built-in tools and documented practices for performance testing:

  • Thespian Benchmark CLI Tool: The framework includes a thespian-bench command that runs standardized performance tests for different transport bases. You can use it to compare throughput and latency across transports (e.g., thespian-bench -b multiprocTCPBase -c 10000).
  • Python Profiling Integration: Pair Thespian with Python’s built-in profiling tools like cProfile to pinpoint bottlenecks. Wrap your main() function with profiling to see if time is being eaten up by serialization, transport, or Actor processing logic.
  • Message Queue Monitoring: Use Thespian’s actorSystemUpdate method to inspect Actor message queue lengths. Adding queue-tracking logic to Actor2 can confirm if backlog is the primary source of delay.
  • Official Documentation Guidance: The Thespian docs explicitly call out batch processing and transport selection as key for low-latency use cases. They also note that TCP-based transports are slower than shared-memory queues for single-machine deployments.

内容的提问来源于stack exchange,提问作者jonatelo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 09:59:53