队列服务支持通过HTTP API、SDK访问,用于发送数据、查询推理结果和管理队列。
队列调用地址
异步推理服务部署完成后,会自动生成输入队列和输出队列(sink队列)两类地址,以HTTP接口为例,说明如下:
地址类型 | 地址格式 | 示例 |
输入队列地址 |
|
|
输出队列地址 |
|
|
在推理服务页签,单击目标服务名称进入概览页面,在基本信息区域,单击查看调用信息,查看输入队列地址、输出队列地址和Token。
在调用信息弹窗中,选择共享网关 > 异步调用页签,调用地址分为公网和 VCP 两类,各包含输入调用地址和输出调用地址。
向队列服务发送数据
使用curl命令向输入队列发送一条同步请求或异步推理请求,具体代码示例如下。
curl -v http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' -d '[{}]'示例结果如下:
> POST /api/predict/qservice HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 19
< X-Eas-Queueservice-Request-Id: 4e034bnvb-e783-4272-9333-68x6a1v8dc6x
<
1033其中:
Response Header中返回的X-Eas-Queueservice-Request-Id,为该请求对应的Request ID:4e034bnvb-e783-4272-9333-68x6a1v8dc6x,可以通过该Request ID来查询数据。
Response Body中返回的是当前请求在队列中的Index:1033,您可以通过Index在当前队列中查询数据。
发送优先数据
在队列服务中,普通数据按照FIFO顺序进行推送,但是在很多场景中,部分数据需要被优先推送和处理。队列服务支持数据优先推送。您可以通过增加query参数_priority_=1,向队列服务推送优先数据。
$ curl -v http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice?_priority_=1 -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' -d '[{}]'示例结果如下:
> POST /api/predict/qservice?_priority_=1 HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 19
< X-Eas-Queueservice-Request-Id: 4033eb55-e783-4922-9777-68d6a1383c76
<
1034优先数据一旦被写入队列,将被优先推送给订阅者,从而进行优先处理。
查看队列服务详情
如果您在向队列服务发送请求时,增加_attrs_=true参数,返回结果中会显示当前队列的详情信息。具体代码示例如下。
curl -v -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice?_attrs_=true示例结果如下:
> GET /api/predict/qservice?_attrs_=true HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 320
<
{"consumers.stats.total":"0","consumers.status.total":"0","meta.header.group":"X-EAS-QueueService-Gid","meta.header.priority":"X-EAS-QueueService-Priority","meta.header.user":"X-EAS-QueueService-Uid","stream.maxPayloadBytes":"524288","meta.name":"pmml_test","meta.state":"Normal","stream.approxMaxLength":"4095","stream.firstEntry":"0","stream.lastEntry":"0","stream.length":"1"}上述结果中返回JSON格式的详情信息,其中关键字段说明如下:
字段名 | 描述 |
stream.maxPayloadBytes | 队列中允许的每个数据项的大小上限,单位为Byte。 |
stream.approxMaxLength | 队列中能存储的数据项的数量上限。 |
stream.firstEntry | 队列中第一个数据项的index。 |
stream.lastEntry | 队列中最后一个数据项的index。 |
stream.length | 队列中当前存储的数据项的数量。 |
meta.state | 当前队列的状态。 |
您也可以在模型在线服务(EAS)页面,单击异步推理服务的名称进入概览页面,然后切换至异步队列页签查看队列信息。
在 PAI-EAS 服务详情页中,单击异步队列 Tab,可查看队列的基本信息(所属资源组、创建时间、单一输入请求最大数据及单一输出返回最大数据)、服务部署资源(实例数、CPU、内存),以及输入队列的当前存储数据数量与实例处理状态。
查询数据
根据条件查询结果
当只使用一个队列服务时,您可以通过Index或Request ID从输入队列中查询数据,具体代码示例如下。
# 通过index查询数据。
$ curl -v -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice?_index_=1022
# 通过request id查询数据。
$ curl -v -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice?requestId=87633037-39a4-40bf-8405-14f8e0c31896示例结果如下:
> GET /api/predict/qservice?_index_=1022&_auto_delete_=false HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 4
< Content-Type: text/plain; charset=utf-8
<
[{}]您可以配置以下参数来查询推理结果,具体参数说明如下:
参数 | 类型 | 核心参数说明 |
_index_ | INT | 要查询数据的起始index。默认为0,表示从队列的初始数据项开始查询,该index越接近被查询数据,查询的效率越高。 |
_length_ | INT | 要查询的数据项的条数。默认为1,表示仅查询一条数据项。 |
_auto_delete_ | BOOL | 是否从队列中删除已查询的数据。默认为TRUE,表示查询完成后,将查询出的数据项自动从队列中删除。 |
_timeout_ | STRING | 超时时间。默认为0,表示查询时队列中无符合要求的数据则立即返回204状态码,否则等待指定时间,在超时时间内如果队列中出现符合要求的数据,则将数据返回。示例值:1s(1秒), 1m(1分钟)。 |
requestId | STRING | requestId为内建的tag,表示通过该tag来查询数据。 说明 当使用异步推理服务功能时,请求从输入队列返回,由EAS服务框架读取输出数据进行处理后将结果自动写入到输出队列中,服务框架会通过requestId这个tag将输入数据与输出数据进行关联,通过输入数据的requestId即可在输出队列中查询结果数据。 |
查询异步推理结果
当队列服务有与之搭配的推理服务时,推理服务会自动从输入队列中读取请求数据,进行推理计算后将推理结果写出到输出队列(sink)中。使用以下代码根据Request ID(0337f7a1-a6f6-49a6-8ad7-ff2fd12b****)从输出队列中查询数据。
$ curl -v -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice/sink?requestId=0337f7a1-a6f6-49a6-8ad7-ff2fd12bbe2d示例结果如下:
> GET /api/predict/qservice/sink?requestId=0337f7a1-a6f6-49a6-8ad7-ff2fd12b**** HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 53
< Content-Type: text/plain; charset=utf-8
<
[{"p_0":0.5224580736905329,"p_1":0.4775419263094671}]清理数据
当您的队列中不再需要某些数据时,可以通过API对数据进行清理。数据清理的方式主要有两种,分别是单条数据删除(delete)和数据截止删除(truncate)。
删除单条数据
# 通过index删除数据。
$ curl -XDELETE -v -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice?_index_=1022示例结果如下:
> DELETE /api/predict/qservice?_index_=1022 HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 4
< Content-Type: text/plain; charset=utf-8
<
OK您可以配置以下参数来查询推理结果,具体参数说明如下:
参数 | 类型 | 核心参数说明 |
_index_ | INT | 要删除的数据index。 |
批量数据删除
# 通过index删除数据。
$ curl -XDELETE -v -H 'Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' http://182848887922****.cn-shanghai.pai-eas.aliyuncs.com/api/predict/qservice?_index_=1023&_trunc_=true示例结果如下:
> DELETE /api/predict/qservice?_index_=1023&_trunc_=true HTTP/1.1
> Host: 182848887922****.cn-shanghai.pai-eas.aliyuncs.com
> Authorization: YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==
>
< HTTP/1.1 200 OK
< Content-Length: 4
< Content-Type: text/plain; charset=utf-8
<
OK您可以配置以下参数来查询推理结果,具体参数说明如下:
参数 | 类型 | 核心参数说明 |
_index_ | INT | 要删除的数据截止index,低于(不包含)这个index的数据将被删除。 |
_trunc_ | BOOL | 在批量删除时必须为true,否则将转换为单条删除。 |
队列服务订阅推送
在异步推理场景中,除了上述的阻塞查询,还可以通过订阅的方式来获取推理结果。队列服务提供了订阅(watch)接口,客户端可以通过该接口来获取推理结果。队列服务根据当前推理服务实例配置的并发数(worker_threads)来控制订阅的窗口(Window)大小,当队列中被写入新数据时,队列服务会自动将数据推送给正在订阅的客户端。
该功能在SDK中基于WebSocket协议封装了客户端实现QueueClient,通过长连接的方式建立推送链路。下面以一个典型的视频、语音流处理场景为例,介绍如何通过Python SDK中的QueueClient来订阅队列中的数据。
推理服务不是必需的,您也可以通过SDK在自定义的服务中订阅队列服务的输入队列,输出结果也可以选择写入到第三方的消息队列中或其它目标存储中(比如输出图片到OSS)。
安装EAS Python SDK。
pip install eas_prediction --user通过QueueClient的
put()方法向输入队列中发送数据,并使用watch()方法从输出队列中订阅数据。在实际使用场景中,发送数据和订阅数据可以由不同的线程处理,本示例中发送数据和订阅数据在同一线程中完成,先put数据,后watch结果。
#!/usr/bin/env python from eas_prediction import QueueClient # 将domain、service_name、token替换为服务实际值 domain = '182848887922***.cn-shanghai.pai-eas.aliyuncs.com' service_name = 'qservice' token = 'YmE3NDkyMzdiMzNmMGM3ZmE4ZmNjZDk0M2NiMDA3OTZmNzc1MT****==' # 创建输入队列对象,用于写入输入数据。 input_queue = QueueClient(domain, service_name) # 如果需要自定义user和group,可以分别通过uid和gid进行指定,示例如下: # input_queue = QueueClient(domain, service_name, uid='your_user_id', gid='your_group_id') input_queue.set_token(token) input_queue.init() # 创建输出队列对象,用于订阅读取输出结果数据。 sink_queue = QueueClient(domain, f'{service_name}/sink') sink_queue.set_token(token) sink_queue.init() # 各输入队列中推送10个数据项。 for x in range(10): index, request_id = input_queue.put('[{}]') print(index, request_id) # 查看输入队列的详情。 attrs = input_queue.attributes() print(attrs) # 从输出队列中watch数据,窗口为5。 i = 0 watcher = sink_queue.watch(0, 5, auto_commit=False) for x in watcher.run(): print(x.data.decode('utf-8')) # 每次收到一个请求数据后处理完成后手动commit。 sink_queue.commit(x.index) i += 1 if i == 10: break # 关闭已经打开的watcher对象,每个客户端实例只允许存在一个watcher对象,若watcher对象不关闭,再次运行时会报错。 watcher.close()