在 Twisted 回调中使用 async/await 语法
Using async/await syntax with Twisted callbacks
我想将 async/await
语法与用于此目的的 Twisted Deferred.addCallback method. But as stated in the documentation, addCallback
callback is called synchronously. I've seen inlineCallbacks 装饰器一起使用,但我更喜欢使用 async/await
语法(如果可能或有意义的话) ).
我从 pika documentation 中获取了原始代码,但我没有运气尝试将其迁移到 async/await 语法:
import pika
from pika import exceptions
from pika.adapters import twisted_connection
from twisted.internet import defer, reactor, protocol, task
async def run_async(connection):
channel = await connection.channel()
exchange = await channel.exchange_declare(exchange='topic_link',type='topic')
queue = await channel.queue_declare(queue='hello', auto_delete=False, exclusive=False)
await channel.queue_bind(exchange='topic_link', queue='hello', routing_key='hello.world')
await channel.basic_qos(prefetch_count=1)
queue_object, consumer_tag = await channel.basic_consume(queue='hello', no_ack=False)
l = task.LoopingCall(read_async, queue_object)
l.start(0.01)
async def read_async(queue_object):
ch,method,properties,body = await queue_object.get()
if body:
print(body)
await ch.basic_ack(delivery_tag=method.delivery_tag)
parameters = pika.ConnectionParameters()
cc = protocol.ClientCreator(reactor, twisted_connection.TwistedProtocolConnection, parameters)
d = cc.connectTCP('rabbitmq', 5672)
d.addCallback(lambda protocol: protocol.ready)
d.addCallback(run_async)
reactor.run()
这显然不起作用,因为没有人等待 run_async
函数。
正如 notorious.no 和 Twisted 文档所指出的那样,ensureDeferred
是正确的选择。不过,您必须包装回调结果,而不是我不清楚的回调本身。
这是最终的样子:
def ensure_deferred(f):
@functools.wraps(f)
def wrapper(*args, **kwargs):
result = f(*args, **kwargs)
return defer.ensureDeferred(result)
return wrapper
@ensure_deferred
async def run(connection):
channel = await connection.channel()
exchange = await channel.exchange_declare(exchange='topic_link', type='topic')
queue = await channel.queue_declare(queue='hello', auto_delete=False, exclusive=False)
await channel.queue_bind(exchange='topic_link', queue='hello', routing_key='hello.world')
await channel.basic_qos(prefetch_count=1)
queue_object, consumer_tag = await channel.basic_consume(queue='hello', no_ack=False)
l = task.LoopingCall(read, queue_object)
l.start(0.01)
@ensure_deferred
async def read(queue_object):
ch, method, properties, body = await queue_object.get()
if body:
print(body)
await ch.basic_ack(delivery_tag=method.delivery_tag)
parameters = pika.ConnectionParameters()
cc = protocol.ClientCreator(reactor, twisted_connection.TwistedProtocolConnection, parameters)
d = cc.connectTCP('rabbitmq', 5672)
d.addCallback(lambda protocol: protocol.ready)
d.addCallback(run)
reactor.run()
谢谢。
我想将 async/await
语法与用于此目的的 Twisted Deferred.addCallback method. But as stated in the documentation, addCallback
callback is called synchronously. I've seen inlineCallbacks 装饰器一起使用,但我更喜欢使用 async/await
语法(如果可能或有意义的话) ).
我从 pika documentation 中获取了原始代码,但我没有运气尝试将其迁移到 async/await 语法:
import pika
from pika import exceptions
from pika.adapters import twisted_connection
from twisted.internet import defer, reactor, protocol, task
async def run_async(connection):
channel = await connection.channel()
exchange = await channel.exchange_declare(exchange='topic_link',type='topic')
queue = await channel.queue_declare(queue='hello', auto_delete=False, exclusive=False)
await channel.queue_bind(exchange='topic_link', queue='hello', routing_key='hello.world')
await channel.basic_qos(prefetch_count=1)
queue_object, consumer_tag = await channel.basic_consume(queue='hello', no_ack=False)
l = task.LoopingCall(read_async, queue_object)
l.start(0.01)
async def read_async(queue_object):
ch,method,properties,body = await queue_object.get()
if body:
print(body)
await ch.basic_ack(delivery_tag=method.delivery_tag)
parameters = pika.ConnectionParameters()
cc = protocol.ClientCreator(reactor, twisted_connection.TwistedProtocolConnection, parameters)
d = cc.connectTCP('rabbitmq', 5672)
d.addCallback(lambda protocol: protocol.ready)
d.addCallback(run_async)
reactor.run()
这显然不起作用,因为没有人等待 run_async
函数。
正如 notorious.no 和 Twisted 文档所指出的那样,ensureDeferred
是正确的选择。不过,您必须包装回调结果,而不是我不清楚的回调本身。
这是最终的样子:
def ensure_deferred(f):
@functools.wraps(f)
def wrapper(*args, **kwargs):
result = f(*args, **kwargs)
return defer.ensureDeferred(result)
return wrapper
@ensure_deferred
async def run(connection):
channel = await connection.channel()
exchange = await channel.exchange_declare(exchange='topic_link', type='topic')
queue = await channel.queue_declare(queue='hello', auto_delete=False, exclusive=False)
await channel.queue_bind(exchange='topic_link', queue='hello', routing_key='hello.world')
await channel.basic_qos(prefetch_count=1)
queue_object, consumer_tag = await channel.basic_consume(queue='hello', no_ack=False)
l = task.LoopingCall(read, queue_object)
l.start(0.01)
@ensure_deferred
async def read(queue_object):
ch, method, properties, body = await queue_object.get()
if body:
print(body)
await ch.basic_ack(delivery_tag=method.delivery_tag)
parameters = pika.ConnectionParameters()
cc = protocol.ClientCreator(reactor, twisted_connection.TwistedProtocolConnection, parameters)
d = cc.connectTCP('rabbitmq', 5672)
d.addCallback(lambda protocol: protocol.ready)
d.addCallback(run)
reactor.run()
谢谢。