MCPcopy Create free account
hub / github.com/PaddlePaddle/FastDeploy / Queue

Class Queue

fastdeploy/inter_communicator/fmq.py:196–268  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

194
195
196class Queue(BaseComponent):
197 def __init__(self, context, name: str, role: str = "producer"):
198 endpoint = EndpointManager.get_endpoint(name)
199 super().__init__(context, endpoint)
200
201 self.name = name
202 self.role = Role(role)
203 self.copy = endpoint.copy
204 self.socket_conf = EndpointManager.config.socket_config
205 self._msg_id = 0
206
207 full_ep = f"{endpoint.protocol}://{endpoint.address}"
208
209 self.socket = self.context.socket(zmq.PUSH if self.role == Role.PRODUCER else zmq.PULL)
210 self.socket_conf.apply(self.socket, self.role == Role.PRODUCER)
211
212 if self.role == Role.PRODUCER:
213 self.socket.connect(full_ep)
214 else:
215 self.socket.bind(full_ep)
216
217 fmq_logger.info(f"Queue {name}({role}) initialized on {full_ep}")
218
219 async def put(self, data: Any, shm_threshold: int = 1024 * 1024):
220 """
221 Send data to the queue.
222
223 Args:
224 data: The data to send. Can be any serializable object or bytes.
225 shm_threshold: Size threshold in bytes. If the data is of type bytes and its size is
226 greater than or equal to this threshold, shared memory will be used to send the message.
227 Default is 1MB (1024 * 1024 bytes).
228
229 Raises:
230 PermissionError: If called by a non-producer role.
231 """
232 if self.role != Role.PRODUCER:
233 raise PermissionError("Only producers can send messages.")
234
235 desc = None
236 payload = data
237
238 if isinstance(data, bytes) and len(data) >= shm_threshold:
239 desc = Descriptor.create(data)
240 payload = None
241
242 msg = Message(msg_id=self._msg_id, payload=payload, descriptor=desc)
243 raw = msg.serialize()
244
245 async with self.lock:
246 await self.socket.send(raw, copy=self.copy)
247 self._msg_id += 1
248
249 async def get(self, timeout: int = None) -> Optional[Message]:
250 # Receive data from queue
251 if self.role != Role.CONSUMER:
252 raise PermissionError("Only consumers can get messages.")
253

Callers 6

pre_compile_from_configFunction · 0.85
__init__Method · 0.85
queueMethod · 0.85
run_with_timeoutFunction · 0.85

Calls

no outgoing calls

Tested by 2