Thespian框架Actor间低延迟消息传输性能低下问题排查
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:
- Is there a problem or oversight in my code?
- Are there faster ways to send messages between Actors?
- 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
datadictionary 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:
multiprocTCPBaseuses 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
multiprocQueueBaseinstead ofmultiprocTCPBase. It uses in-memory queues (no TCP serialization/network overhead) which is way faster for single-host inter-Actor communication. For cross-machine needs, considermultiprocUDPBase(just note UDP’s reliability tradeoffs for low-latency use cases). - Trim Unnecessary Message Data: You’re sending the entire
datadict (includinglimandtic) every time. Instead, send only the critical payload, and pass metadata likelimonce 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-benchcommand 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
cProfileto pinpoint bottlenecks. Wrap yourmain()function with profiling to see if time is being eaten up by serialization, transport, or Actor processing logic. - Message Queue Monitoring: Use Thespian’s
actorSystemUpdatemethod 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

